diff --git a/docs/PRD_ADMIN_MODERATION.md b/docs/PRD_ADMIN_MODERATION.md index 032c2aa..b17fe10 100644 --- a/docs/PRD_ADMIN_MODERATION.md +++ b/docs/PRD_ADMIN_MODERATION.md @@ -41,7 +41,7 @@ This is label-integrated moderation, rather than a local deletion feature with f - **Labels for public subjects are published by default** (2026-09-20); restricted subjects are excluded until sharing rules exist. - **The labeling authority is the instance DID** (2026-09-20, `INSTANCE_DID`, a `did:web`), with a separate `#atproto_label` key and an `#atproto_labeler` service added to its DID document. No dedicated labeler account. Whether to reuse an existing label publisher remains open. - **Trust configuration is operator-managed** (2026-09-20), with attributed changes and no per-user opt-out of mandatory removals. -- **Removed posts use `#moderatedPost` wherever the schema allows it** (2026-09-20); `#notFoundPost` is the fallback only for an endpoint where a supported older decoder fails on the new variant. +- **Removed posts use `#moderatedPost` wherever the schema allows it** (2026-09-20). The `#notFoundPost` fallback this decision allowed was not built: if a supported older decoder fails on the new variant, admin actions stay disabled (`MODERATION_ADMINS` empty) until supported builds pass, as in `docs/LEXICON_PUBLISHING.md` release gate 5. - **Public reason codes** (2026-09-20): `spam`, `harassment`, `doxing`, `illegal-content`, `rule-violation`, `moderator-discretion`. A private `csam` classification is shown publicly as `illegal-content`. Public modlog entries with reason `illegal-content` or `doxing` omit the subject URI, because the record is still fetchable from its author's PDS and the log must not work as a link to it. Decided 2026-09-29 (Q-ORACLE): omitting the subject keeps doxing and illegal-content subjects out of the browsable and filterable log, but an observer who already holds a URI can see that it is removed and, finding no public entry for it, infer that the reason is doxing or illegal content; the user accepted this on 2026-09-29. - **No public free-text explanation in this release** (2026-09-20). The public log carries the reason code only; private notes stay admin-only. With no public text there is no correction procedure to define. - **No write-path refusal for removed subjects** (2026-09-20). Clients hide controls on placeholders; the consumer indexes normally (section 5.2). @@ -525,7 +525,7 @@ The published `comment.defs#commentView` requires a verbatim `record`, but the A | `community/comment/defs.json#blockedComment` | Description-only: `blockedBy: moderator` is published but never emitted; mark it unused. Moderator removal is represented only by the `commentView` placeholder. | Do not remove the known value (breaking) and do not start emitting it as a second representation of removal. | | `community/comment/getComments.json` | No schema change. A removed post returns the existing `NotFound` error (decided 2026-09-20). Removed comments appear as placeholders at any depth with their visible descendants intact. | One endpoint for old and new readers. Preserved discussion under a removed post is later scope and would need a versioned endpoint with a post-header union. | | `internal/core/comments/view_models.go` | Correct the stale doc comments naming `getComments#commentView` and `social.coves.feed.getComments`; the NSIDs are `social.coves.community.comment.defs#commentView` and `social.coves.community.comment.getComments`. | Comment-only. | -| `community/post/defs.json` | Add `#moderatedPost`: required `uri` and `moderation` (`moderation.defs#moderationView`, semantically removed); optional `authorDid` (`did`) and safe community reference. No record/title/embed. Retain `#removedPost` for its existing community-removal meaning. Update `#blockedPost` and its `blockedBy` descriptions so their reference to `removedPost` means community removals; instance moderation uses `moderatedPost` or the documented unavailable fallback. | New object instead of relabeling a community's decision as an instance decision. Source attribution distinguishes local/inherited restrictions without pretending to be a human actor. The blocked-view wording update is description-only and compatible. | +| `community/post/defs.json` | Add `#moderatedPost`: required `uri` and `moderation` (`moderation.defs#moderationView`, semantically removed); optional `authorDid` (`did`) and safe community reference. No record/title/embed. Retain `#removedPost` for its existing community-removal meaning. Update `#blockedPost` and its `blockedBy` descriptions so their reference to `removedPost` means community removals; instance moderation uses `moderatedPost`. | New object instead of relabeling a community's decision as an instance decision. Source attribution distinguishes local/inherited restrictions without pretending to be a human actor. The blocked-view wording update is description-only and compatible. | | `community/post/defs.json#postView` | Add optional `moderation` ref to `moderation.defs#moderationView`. Normal NSFW posts keep their verbatim record and ordinary view; active moderator `contentLabels` drive existing sensitivity presentation alongside record self-labels. | Additive optional property, no versioned endpoint required for this metadata. Do not rewrite `record.labels` or return a removal tombstone for a post that is only NSFW. Update client bindings and sensitivity checks before enabling the admin workflow. | | `community/post/get.json` | Add `#moderatedPost` to the existing open `posts.items` union and serve it for instance removals (decided 2026-09-20), subject to the old-reader check below. | Schema-additive, so schema evolution alone does not require `post.getV2`; runtime client compatibility must still be demonstrated. Unknown variants must render unavailable rather than crash a batch or expose raw JSON. | | `actor/getComments.json` | Keep its normal `commentView` response. Filter moderated comments, comments under a removed post, and any author-deleted placeholders from profile activity before hydration. | No schema change. Audit whether placeholders are currently returned here; profile tombstones have no use. | @@ -535,7 +535,7 @@ The published `comment.defs#commentView` requires a verbatim `record`, but the A **Old-reader contract:** comments need no separate old-reader path, because removed comments use the wire shape installed clients already decode and a removed post's thread uses the already-declared `NotFound`. Before enabling actions, pin with tests that supported web and mobile builds render `deletionReason: moderator`, tolerate the optional `moderation` object, and handle thread `NotFound` without a crash loop. Every other still-served endpoint must redact using only a response its published schema allows. -For **every unversioned endpoint gaining a union variant**, test actual supported older mobile and web decoders, not just newly generated clients. For `post.get`, include normal posts plus an unknown variant in a 25-URI response and verify one removed post cannot fail the entire batch. If any supported decoder fails, keep instance removals as the existing `#notFoundPost` for all callers of that endpoint until compatibility is established; never substitute the community-only `#removedPost`. If incompatible readers must remain supported while new clients need source-attributed post tombstones, select and inventory a versioned post-read endpoint before rollout. Do not assume installed apps upgrade or infer decoder support from authentication. Never serve a normal `commentView` containing removed content. +For **every unversioned endpoint gaining a union variant**, test actual supported older mobile and web decoders, not just newly generated clients. For `post.get`, include normal posts plus an unknown variant in a 25-URI response and verify one removed post cannot fail the entire batch. Until supported client builds pass, leave `MODERATION_ADMINS` empty so no admin action can create a removal (`docs/LEXICON_PUBLISHING.md` release gate 5): `post.get` serves `#moderatedPost` for every instance removal and has no `#notFoundPost` fallback switch. Never substitute the community-only `#removedPost`. If incompatible readers must remain supported while new clients need source-attributed post tombstones, select and inventory a versioned post-read endpoint before rollout. Do not assume installed apps upgrade or infer decoder support from authentication. Never serve a normal `commentView` containing removed content. ### 14.5 Errors, privacy, and unknown values diff --git a/internal/atproto/jetstream/authorpost.go b/internal/atproto/jetstream/authorpost.go index ed40a6f..5dd42ac 100644 --- a/internal/atproto/jetstream/authorpost.go +++ b/internal/atproto/jetstream/authorpost.go @@ -720,7 +720,7 @@ func (c *PostEventConsumer) upsertAuthorPost(ctx context.Context, authorDID stri PostV2Collection, uri, stored.communityDID, record.Community) // The author's repo still serves the ignored record's images, so a // removed post blocks them for the author (PRD Q-I6). - if err := c.blockIgnoredUpdateMedia(ctx, uri, authorDID, commit.Rev, timeUS, stored.indexedAt, record.Embed); err != nil { + if err := c.blockIgnoredUpdateMedia(ctx, uri, authorDID, commit.Rev, timeUS, record.Embed); err != nil { return err } return nil @@ -839,25 +839,15 @@ func (c *PostEventConsumer) upsertAuthorPost(ctx context.Context, authorDID stri // blockIgnoredUpdateMedia blocks, for the author only, the images of a postv2 // update the consumer ignores, when the post has an active removal. A stale // event (older by rev or by event time than the indexed state) blocks nothing: -// the newer state supersedes it. Neither guard advances anything, because the -// ignored content is never applied. -func (c *PostEventConsumer) blockIgnoredUpdateMedia(ctx context.Context, uri, authorDID, rev string, timeUS int64, storedIndexedAt time.Time, embed map[string]interface{}) error { - if evTime, ok := eventTime(timeUS); ok && !storedIndexedAt.Before(evTime) { - return nil - } - stale, err := recordRevIsStale(ctx, c.db, uri, rev) - if err != nil { - return fmt.Errorf("failed to check rev of ignored post update: %w", err) - } - if stale { - logSkippedStaleRev(ConsumerPosts, "update", uri, rev) - return nil - } +// the newer state supersedes it. Both guards are decided under the gate and +// post row locks, in the transaction that inserts the blocks, and neither +// advances anything, because the ignored content is never applied. +func (c *PostEventConsumer) blockIgnoredUpdateMedia(ctx context.Context, uri, authorDID, rev string, timeUS int64, embed map[string]interface{}) error { _, embedJSON, _, err := serializePostContent(nil, embed, nil) if err != nil { return err } - if err := c.blockIncomingMedia(ctx, uri, authorDID, embedJSON); err != nil { + if err := c.blockIncomingMedia(ctx, uri, authorDID, rev, timeUS, embedJSON); err != nil { return fmt.Errorf("failed to block media of ignored post update: %w", err) } return nil diff --git a/internal/atproto/jetstream/comment_consumer.go b/internal/atproto/jetstream/comment_consumer.go index 731db25..d46b726 100644 --- a/internal/atproto/jetstream/comment_consumer.go +++ b/internal/atproto/jetstream/comment_consumer.go @@ -360,9 +360,9 @@ func (c *CommentEventConsumer) updateComment(ctx context.Context, repoDID string log.Printf(" Existing parent: %s (CID: %s)", storedParentURI, storedParentCID) log.Printf(" Incoming parent: %s (CID: %s)", commentRecord.Reply.Parent.URI, commentRecord.Reply.Parent.CID) // The commenter's repo still serves the rejected record's images, so a - // removed comment blocks them for the commenter (PRD Q-I6). The recency - // guard above already ran; a rev-stale event blocks nothing either. - if err := c.blockRejectedUpdateMedia(ctx, uri, repoDID, commit.Rev, commentRecord.Embed); err != nil { + // removed comment blocks them for the commenter (PRD Q-I6). A stale + // event, by rev or by event time, blocks nothing. + if err := c.blockRejectedUpdateMedia(ctx, uri, repoDID, commit.Rev, timeUS, commentRecord.Embed); err != nil { return err } // PERMANENT: threading reassignment is a policy rejection that no retry or @@ -505,22 +505,15 @@ func (c *CommentEventConsumer) updateComment(ctx context.Context, repoDID string } // blockRejectedUpdateMedia blocks the images of a rejected comment update for -// the commenter when the comment has an active removal. A rev-stale event -// blocks nothing, and the read-only rev check advances nothing, because the -// rejected content is never applied. -func (c *CommentEventConsumer) blockRejectedUpdateMedia(ctx context.Context, uri, commenterDID, rev string, embed map[string]interface{}) error { +// the commenter when the comment has an active removal. An event stale by rev +// or by event time blocks nothing. Both guards are decided under the gate and +// comment row locks, in the transaction that inserts the blocks, and advance +// nothing, because the rejected content is never applied. +func (c *CommentEventConsumer) blockRejectedUpdateMedia(ctx context.Context, uri, commenterDID, rev string, timeUS int64, embed map[string]interface{}) error { blobCIDs := embeds.CommentImageCIDs(embed) if c.mediaReconciler == nil || len(blobCIDs) == 0 { return nil } - stale, err := recordRevIsStale(ctx, c.db, uri, rev) - if err != nil { - return fmt.Errorf("failed to check rev of rejected comment update: %w", err) - } - if stale { - logSkippedStaleRev(ConsumerComments, "update", uri, rev) - return nil - } tx, err := c.db.BeginTx(ctx, nil) if err != nil { return fmt.Errorf("failed to begin media block transaction: %w", err) @@ -530,6 +523,13 @@ func (c *CommentEventConsumer) blockRejectedUpdateMedia(ctx context.Context, uri log.Printf("Failed to rollback transaction: %v", rollbackErr) } }() + current, err := lockCurrentIncomingEvent(ctx, tx, lockCommentRowQuery, ConsumerComments, uri, rev, timeUS) + if err != nil { + return fmt.Errorf("failed to check rejected comment update is current: %w", err) + } + if !current { + return nil + } if err := c.commitIncomingCommentMedia(ctx, tx, uri, commenterDID, blobCIDs); err != nil { return fmt.Errorf("failed to block media of rejected comment update: %w", err) } @@ -723,6 +723,12 @@ func (c *CommentEventConsumer) indexCommentAndUpdateCounts(ctx context.Context, log.Printf("Comment already indexed: %s (idempotent replay)", comment.URI) // The dropped record's images are still served from the commenter's // repo, so a removed comment blocks them for the commenter (PRD Q-I6). + // The gate claim above already orders this against same-record + // events; the row lock orders it against a removal that has read + // the comment, as the re-create UPDATE above would. + if _, _, lockErr := lockContentRow(ctx, tx, lockCommentRowQuery, comment.URI); lockErr != nil { + return lockErr + } if commitErr := c.commitIncomingCommentMedia(ctx, tx, comment.URI, comment.CommenterDID, incomingCommentImageCIDs(comment.Embed)); commitErr != nil { return fmt.Errorf("failed to commit transaction: %w", commitErr) diff --git a/internal/atproto/jetstream/moderation_rejected_media_test.go b/internal/atproto/jetstream/moderation_rejected_media_test.go index 7c362c9..ea52385 100644 --- a/internal/atproto/jetstream/moderation_rejected_media_test.go +++ b/internal/atproto/jetstream/moderation_rejected_media_test.go @@ -3,6 +3,9 @@ package jetstream import ( + "context" + "errors" + "fmt" "net/http" "testing" "time" @@ -192,3 +195,98 @@ func TestModerationStaleRejectedEventAddsNoBlocks(t *testing.T) { c.assertImageBNotBlocked(t) }) } + +// holdNewerEvent writes, in a transaction it leaves open until commit is +// called, what a newer same-record event writes: its rev gate claim, the +// first statement of every consumer write, and the row's indexed_at watermark. +func (c *rejectedMediaCase) holdNewerEvent(t *testing.T, rev string, at time.Time) (commit func()) { + t.Helper() + table := "posts" + if c.subjectIsComment { + table = "comments" + } + tx, err := c.h.db.BeginTx(t.Context(), nil) + require.NoError(t, err) + t.Cleanup(func() { _ = tx.Rollback() }) + won, err := tryAdvanceRecordRev(t.Context(), tx, c.subject.URI, rev) + require.NoError(t, err) + require.True(t, won, "the held event must be newer than the indexed state") + _, err = tx.ExecContext(t.Context(), `UPDATE `+table+` SET indexed_at = $2 WHERE uri = $1`, c.subject.URI, at) + require.NoError(t, err) + return func() { require.NoError(t, tx.Commit()) } +} + +// A refused event must decide that it is current in the transaction that +// inserts its blocks. A newer event that commits after the refused event's +// first read of the indexed state, but before its reconciliation, supersedes +// it, and the refused event must then add no blocks. +func TestModerationRejectedEventRechecksFreshnessUnderLock(t *testing.T) { + type deliver func(ctx context.Context, c *rejectedMediaCase, rev string, at time.Time) error + retarget := func(ctx context.Context, c *rejectedMediaCase, rev string, at time.Time) error { + return c.h.posts.HandleEvent(ctx, c.postCommunityChange(rev, at)) + } + rethread := func(ctx context.Context, c *rejectedMediaCase, rev string, at time.Time) error { + err := c.h.comments.HandleEvent(ctx, editMediaCommentEvent(c.otherPost, "update", c.commentRkey, rev, + editMediaCID("rethreaded comment record"), at, c.imageA, c.imageB)) + if errors.Is(err, ErrPermanentEvent) { + return nil + } + return fmt.Errorf("the threading change must still be rejected, got %v", err) + } + for _, path := range []struct { + name string + comment bool + deliver deliver + }{ + {name: "postv2 community change", deliver: retarget}, + {name: "comment threading change", comment: true, deliver: rethread}, + } { + for _, newer := range []struct { + name string + // newerRev and newerAt are the held event's; refusedRev and + // refusedAt the refused event's, as indexes into revs and seconds + // after the case started. + newerRev, refusedRev int + newerAt, refusedAt time.Duration + }{ + // A cross-feed copy carries a newer emission time, so only rev orders it. + {name: "newer by rev", newerRev: 2, refusedRev: 1, newerAt: 2 * time.Second, refusedAt: 3 * time.Second}, + {name: "newer by event time", newerRev: 1, refusedRev: 2, newerAt: 2 * time.Second, refusedAt: time.Second}, + } { + t.Run(path.name+" superseded "+newer.name, func(t *testing.T) { + c := newRejectedMediaCase(t, "fresh "+path.name+" "+newer.name, path.comment, true) + // The pool is three connections: the held event, the refused + // event and this one. + observer, err := c.h.db.Conn(t.Context()) + require.NoError(t, err) + t.Cleanup(func() { _ = observer.Close() }) + commit := c.holdNewerEvent(t, c.revs[newer.newerRev], c.started.Add(newer.newerAt)) + + ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) + defer cancel() + refusedDone := make(chan error, 1) + go func() { + refusedDone <- path.deliver(ctx, c, c.revs[newer.refusedRev], c.started.Add(newer.refusedAt)) + }() + testkit.WaitFor(t, 5*time.Second, func() (bool, error) { + if len(refusedDone) > 0 { + return true, nil + } + var waiting int + err := observer.QueryRowContext(t.Context(), ` + SELECT count(*) FROM pg_stat_activity + WHERE datname = current_database() AND pid <> pg_backend_pid() + AND wait_event_type = 'Lock' AND wait_event <> 'advisory' + `).Scan(&waiting) + return waiting > 0, err + }, testkit.WithDescription("refused event waiting on a lock or finished")) + assert.Empty(t, refusedDone, "the refused event must not reconcile while a newer event is uncommitted") + + commit() + require.NoError(t, <-refusedDone) + c.assertContentUnchanged(t) + c.assertImageBNotBlocked(t) + }) + } + } +} diff --git a/internal/atproto/jetstream/moderation_removal_edit_lock_integration_test.go b/internal/atproto/jetstream/moderation_removal_edit_lock_integration_test.go index c20b80d..a028a1e 100644 --- a/internal/atproto/jetstream/moderation_removal_edit_lock_integration_test.go +++ b/internal/atproto/jetstream/moderation_removal_edit_lock_integration_test.go @@ -47,12 +47,35 @@ func (tx *pausingModerationTransaction) InsertAction(ctx context.Context, action return tx.Transaction.InsertAction(ctx, action) } +// rejectedLockCommentConsumer indexes a comment on the fixture's post with +// one image, and another post the rejected events re-parent it to. +func rejectedLockCommentConsumer(t *testing.T, f postModerationConsumerFixture, rev string, now time.Time) (*CommentEventConsumer, moderation.StrongRef, moderation.StrongRef, string) { + t.Helper() + reconciler := moderation.NewMediaReconciler(postgres.NewModerationRepository(f.db), fixtures.InstanceDID(), nil) + consumer := NewCommentEventConsumer(postgres.NewCommentRepository(f.db), f.db, WithCommentMediaReconciler(reconciler)) + otherRkey := testkit.TID() + otherPost := moderation.StrongRef{URI: pv2URI(pv2Author, otherRkey), CID: editMediaCID("locked other post " + otherRkey)} + require.NoError(t, f.consumer.HandleEvent(t.Context(), pv2Event(pv2Author, "create", otherRkey, f.revs[0], + otherPost.CID, f.createdAt, postModerationRecord()))) + root := moderation.StrongRef{URI: f.uri, CID: postModerationRecordCID} + rkey := testkit.TID() + created := editMediaCID("locked comment create " + rkey) + require.NoError(t, consumer.HandleEvent(t.Context(), + editMediaCommentEvent(root, "create", rkey, rev, created, now, postModerationCIDOne))) + subject := moderation.StrongRef{URI: "at://" + pv2Author + "/" + moderation.CommentCollection + "/" + rkey, CID: created} + return consumer, subject, otherPost, rkey +} + // An author edit that adds image Y must not commit between the removal's read // of the indexed subject and the removal's commit. If it could, the edit's // media reconciliation would find no active removal yet, and the removal would // block only the images it read, leaving Y served. The content-row share lock // taken with that read orders the two: the edit waits for the removal, then its // reconciliation sees the active removal and blocks Y. +// +// The same holds for the events the consumers refuse to index but whose images +// they block (PRD Q-I6): each takes a conflicting lock on the content row +// before it reads the active removal. func TestModerationRemovalHoldsContentRowAgainstAuthorEdit(t *testing.T) { imageY := postModerationCIDTwo for _, scenario := range []struct { @@ -60,6 +83,11 @@ func TestModerationRemovalHoldsContentRowAgainstAuthorEdit(t *testing.T) { // setup indexes the subject with one image and returns it with the // consumer and event that edit it to add imageY. setup func(t *testing.T, f postModerationConsumerFixture) (moderation.StrongRef, func(context.Context) error) + // rejected is true when the consumer refuses to index the event, so + // the subject keeps its indexed record. + rejected bool + // editErr is the error the event returns once the removal commits. + editErr error }{ {name: "postv2", setup: func(t *testing.T, f postModerationConsumerFixture) (moderation.StrongRef, func(context.Context) error) { update := pv2Event(pv2Author, "update", f.rkey, f.revs[1], editMediaCID("locked post edit"), @@ -81,6 +109,28 @@ func TestModerationRemovalHoldsContentRowAgainstAuthorEdit(t *testing.T) { return moderation.StrongRef{URI: "at://" + pv2Author + "/" + moderation.CommentCollection + "/" + rkey, CID: created}, func(ctx context.Context) error { return consumer.HandleEvent(ctx, update) } }}, + {name: "postv2 update that changes community", rejected: true, setup: func(t *testing.T, f postModerationConsumerFixture) (moderation.StrongRef, func(context.Context) error) { + record := postModerationRecord(postModerationCIDOne, imageY) + record["community"] = pv2Prefix + "community2" + update := pv2Event(pv2Author, "update", f.rkey, f.revs[1], editMediaCID("locked post retarget"), + f.createdAt+1_000_000, record) + return moderation.StrongRef{URI: f.uri, CID: postModerationRecordCID}, + func(ctx context.Context) error { return f.consumer.HandleEvent(ctx, update) } + }}, + {name: "comment update that changes threading references", rejected: true, editErr: ErrPermanentEvent, setup: func(t *testing.T, f postModerationConsumerFixture) (moderation.StrongRef, func(context.Context) error) { + now := time.Now() + consumer, subject, otherPost, rkey := rejectedLockCommentConsumer(t, f, f.revs[2], now) + update := editMediaCommentEvent(otherPost, "update", rkey, f.revs[3], editMediaCID("locked comment rethread"), + now.Add(time.Second), postModerationCIDOne, imageY) + return subject, func(ctx context.Context) error { return consumer.HandleEvent(ctx, update) } + }}, + {name: "same-rkey comment re-create with a changed parent", rejected: true, setup: func(t *testing.T, f postModerationConsumerFixture) (moderation.StrongRef, func(context.Context) error) { + now := time.Now() + consumer, subject, otherPost, rkey := rejectedLockCommentConsumer(t, f, f.revs[2], now) + recreate := editMediaCommentEvent(otherPost, "create", rkey, f.revs[3], editMediaCID("locked comment recreate"), + now.Add(time.Second), imageY) + return subject, func(ctx context.Context) error { return consumer.HandleEvent(ctx, recreate) } + }}, } { t.Run(scenario.name, func(t *testing.T) { f := newPostModerationConsumerFixture(t) @@ -152,14 +202,27 @@ func TestModerationRemovalHoldsContentRowAgainstAuthorEdit(t *testing.T) { releaseRemoval() require.NoError(t, <-removalDone) - require.NoError(t, <-editDone) + if scenario.editErr != nil { + require.ErrorIs(t, <-editDone, scenario.editErr) + } else { + require.NoError(t, <-editDone) + } require.NoError(t, observer.QueryRowContext(t.Context(), ` SELECT cid FROM posts WHERE uri = $1 UNION ALL SELECT cid FROM comments WHERE uri = $1 `, subject.URI).Scan(&indexedCID)) - assert.NotEqual(t, subject.CID, indexedCID, "the edit must be indexed after the removal commits") + if scenario.rejected { + assert.Equal(t, subject.CID, indexedCID, "the refused event must not be indexed") + } else { + assert.NotEqual(t, subject.CID, indexedCID, "the edit must be indexed after the removal commits") + } blocked, err := postgres.NewModerationRepository(f.db).IsBlocked(t.Context(), pv2Author, imageY) require.NoError(t, err) assert.True(t, blocked, "the image the edit added must be blocked under the removal") + var ownerless int + require.NoError(t, observer.QueryRowContext(t.Context(), ` + SELECT count(*) FROM moderation_media_blocks WHERE blob_cid = $1 AND owner_did IS NULL + `, imageY).Scan(&ownerless)) + assert.Zero(t, ownerless, "an image added after the removal read the subject must not get an every-owner block") }) } } diff --git a/internal/atproto/jetstream/post_consumer.go b/internal/atproto/jetstream/post_consumer.go index 7714b98..7deacfa 100644 --- a/internal/atproto/jetstream/post_consumer.go +++ b/internal/atproto/jetstream/post_consumer.go @@ -1,6 +1,14 @@ package jetstream import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "log" + "time" + "Coves/internal/atproto/identity" "Coves/internal/core/bridgedvotes" "Coves/internal/core/communities" @@ -9,13 +17,6 @@ import ( "Coves/internal/core/posts" "Coves/internal/core/richtext" "Coves/internal/core/users" - "context" - "database/sql" - "encoding/json" - "errors" - "fmt" - "log" - "time" ) // PostEventConsumer consumes author-owned posts and community decisions from @@ -284,19 +285,12 @@ func (c *PostEventConsumer) applyPostContentUpdate(ctx context.Context, in postC // Skip soft-deleted rows: a deleted post should not be resurrected by an edit. // The author's repo still serves the incoming blobs, so a removed post's // recreated images are blocked even though the content is not indexed. The - // rev gate runs first (read-only: this path never advances the rev) so an - // out-of-order event that predates the delete blocks nothing. + // rev gate is checked under its lock (read-only: this path never advances + // the rev) so an out-of-order event that predates the delete blocks nothing. + // Event time is not compared here (timeUS 0), as this skip never did. if in.storedDeletedAt != nil { - stale, err := recordRevIsStale(ctx, c.db, in.uri, in.rev) - if err != nil { - return false, fmt.Errorf("failed to check rev of soft-deleted post update: %w", err) - } - if stale { - logSkippedStaleRev(ConsumerPosts, "update", in.uri, in.rev) - return false, nil - } log.Printf("Update event for soft-deleted post: %s (skipping)", in.uri) - if err := c.blockIncomingMedia(ctx, in.uri, in.authorDID, in.embed); err != nil { + if err := c.blockIncomingMedia(ctx, in.uri, in.authorDID, in.rev, 0, in.embed); err != nil { return false, fmt.Errorf("failed to block media of skipped post update: %w", err) } return false, nil @@ -722,8 +716,9 @@ func (c *PostEventConsumer) commitIncomingMediaWrite(ctx context.Context, tx *sq } // blockIncomingMedia blocks the blobs of an incoming post event that writes no -// post row. It is a no-op unless the post has an active removal. -func (c *PostEventConsumer) blockIncomingMedia(ctx context.Context, uri, ownerDID string, embed sql.NullString) error { +// post row. It is a no-op unless the post has an active removal, and a no-op +// for an event that is stale by rev, or by event time when timeUS is positive. +func (c *PostEventConsumer) blockIncomingMedia(ctx context.Context, uri, ownerDID, rev string, timeUS int64, embed sql.NullString) error { blobCIDs := incomingPostBlobCIDs(embed) if c.mediaReconciler == nil || len(blobCIDs) == 0 { return nil @@ -737,5 +732,60 @@ func (c *PostEventConsumer) blockIncomingMedia(ctx context.Context, uri, ownerDI log.Printf("Failed to rollback transaction: %v", rollbackErr) } }() + current, err := lockCurrentIncomingEvent(ctx, tx, lockPostRowQuery, ConsumerPosts, uri, rev, timeUS) + if err != nil || !current { + return err + } return c.commitIncomingMediaWrite(ctx, tx, uri, ownerDID, blobCIDs) } + +// The content-row locks taken before an incoming-media active-removal read. +// FOR NO KEY UPDATE conflicts with the FOR SHARE a removal takes when it reads +// the indexed subject, as an ordinary update's row write does. +const ( + lockPostRowQuery = `SELECT indexed_at FROM posts WHERE uri = $1 FOR NO KEY UPDATE` + lockCommentRowQuery = `SELECT indexed_at FROM comments WHERE uri = $1 FOR NO KEY UPDATE` +) + +// lockContentRow locks the indexed row of a record whose incoming media is +// about to be reconciled, so a removal that has already read the subject +// commits first and the active-removal read that follows sees it. found is +// false when the row is gone; no removal can then be holding it. +func lockContentRow(ctx context.Context, tx *sql.Tx, lockQuery, uri string) (indexedAt time.Time, found bool, err error) { + err = tx.QueryRowContext(ctx, lockQuery, uri).Scan(&indexedAt) + if errors.Is(err, sql.ErrNoRows) { + return time.Time{}, false, nil + } + if err != nil { + return time.Time{}, false, fmt.Errorf("failed to lock indexed row %s: %w", uri, err) + } + return indexedAt, true, nil +} + +// lockCurrentIncomingEvent decides inside tx whether an event the consumer +// refuses to index is still the newest for its record, and holds the locks +// that keep that answer and the active-removal read after it true until tx +// ends. It locks in the order every same-record write does: the rev gate row, +// then the content row. A stale event, by rev or by an event time (when timeUS +// is positive) not after the row's indexed_at, returns false. It advances +// nothing, because the refused content is never applied. +func lockCurrentIncomingEvent(ctx context.Context, tx *sql.Tx, lockQuery, consumer, uri, rev string, timeUS int64) (bool, error) { + stale, err := lockedRecordRevIsStale(ctx, tx, uri, rev) + if err != nil { + return false, err + } + if stale { + logSkippedStaleRev(consumer, "update", uri, rev) + return false, nil + } + indexedAt, found, err := lockContentRow(ctx, tx, lockQuery, uri) + if err != nil { + return false, err + } + if evTime, ok := eventTime(timeUS); ok && found && !indexedAt.Before(evTime) { + log.Printf("INFO: not blocking media of stale %s event for %s (event time %s <= last indexed %s)", + consumer, uri, evTime.Format(time.RFC3339Nano), indexedAt.Format(time.RFC3339Nano)) + return false, nil + } + return true, nil +} diff --git a/internal/atproto/jetstream/rev_gate.go b/internal/atproto/jetstream/rev_gate.go index 10022f6..5637e7e 100644 --- a/internal/atproto/jetstream/rev_gate.go +++ b/internal/atproto/jetstream/rev_gate.go @@ -100,6 +100,28 @@ func recordRevIsStale(ctx context.Context, q revGateQuerier, uri, rev string) (b return stale, nil } +// lockedRecordRevIsStale is recordRevIsStale with the gate row locked until tx +// ends. A same-record event's tryAdvanceRecordRev waits on that lock, so no +// newer event can commit between this answer and tx's commit. It advances +// nothing. With no gate row there is nothing to lock and the event is not +// stale. +func lockedRecordRevIsStale(ctx context.Context, tx *sql.Tx, uri, rev string) (bool, error) { + if rev == "" { + return false, nil + } + var stale bool + err := tx.QueryRowContext(ctx, + `SELECT rev >= $2 FROM jetstream_record_revs WHERE record_uri = $1 FOR UPDATE`, uri, rev, + ).Scan(&stale) + if errors.Is(err, sql.ErrNoRows) { + return false, nil + } + if err != nil { + return false, fmt.Errorf("failed to lock record rev for %s: %w", uri, err) + } + return stale, nil +} + // logSkippedStaleRev is the single, grep-able log line for gate skips. // Rejected stale events are the system WORKING (e.g. the bsky feed's delayed // copies of self-feed events); this line makes that observable and diff --git a/internal/db/postgres/moderation_repo.go b/internal/db/postgres/moderation_repo.go index 8f6511f..e48dbd5 100644 --- a/internal/db/postgres/moderation_repo.go +++ b/internal/db/postgres/moderation_repo.go @@ -301,8 +301,8 @@ func (t *moderationTransaction) GetAction(ctx context.Context, actionID string) // committed row, and a restore that arrives later waits at its decision // update for this transaction's blocks to commit, so its later // DeactivateMediaBlocks sees them. Mutations take it after the subject lock -// and the consumer after its own row write; neither then waits on a lock the -// other holds, so the order cannot deadlock. +// and the consumer after its own row write or row lock; neither then waits on +// a lock the other holds, so the order cannot deadlock. func (t *moderationTransaction) ActiveRemoval(ctx context.Context, authorityDID, subjectURI string) (*moderation.Action, error) { action, err := scanModerationAction(t.tx.QueryRowContext(ctx, ` SELECT `+moderationActionColumns+` FROM moderation_actions a