diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index d3e09ee..ddcecbe 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -119,6 +119,46 @@ task documents and git history rather than this list. event. The decision-19 reconciliation job (task 18) is the natural home for a periodic acceptance-vs-record pin audit. +## Community bans (task 17c-3) + +- **The in-transaction ban re-check NARROWS the race, it does not close + it.** `accept()`'s side effect asks `StandingTx` on the acceptance + transaction before writing, so the common ordering is covered — but + under READ COMMITTED a ban committing between that read and the commit + is still possible. Closing it needs SERIALIZABLE or a keyed lock shared + with the ban path, which is a heavier change than the exposure + warrants; the code says so where it happens. +- **A delivery claimed a moment before the ban commits can still be + POSTed.** Cancellation reaches the row, not the socket: the worker's + claim transaction has already committed, and the fencing correctly + discards its settlement afterwards. So the row can read `cancelled` + for something the peer accepted — and the vote reseed subtracts + exactly the `delivered` state. A worker-side ban re-check immediately + before `SendActivityAs` would shrink the window to the send itself; it + was deliberately NOT added, because the escape it was meant to close + (new comments and votes) is now closed at the source, and it would be + defence-in-depth with no test behind it plus a new index. 17e's + reconciliation is the natural place to detect the divergence. +- **A temporary ban that lapses leaves state divergent.** Posts accepted + before the ban keep their acceptance records in the community repo, + but their deliveries were cancelled, and `cancelled` is terminal + (`RedrivePoisoned` revives only poisoned rows). So the community + renders content on Coves that Lemmy will never receive, and nothing + re-queues it when the exclusion ends. Exactly the atproto-vs-outbound + divergence 17e is meant to report. +- **No e2e Block coverage.** The wire keys are hand-written in unit + fixtures, so nothing exercises a real Lemmy's spelling. Both `expires` + and AS2 `endTime` are read, but a third spelling — or a shape + difference in `target` — would be invisible until production. Task 18's + moderation scenario is the place for a captured Block. +- **A ban is recorded only in a private table.** A + `social.coves.moderation.ban` lexicon exists, and every other + moderation decision here is published into the community repo + (acceptance, removal) where Coves' surfaces read it. As landed, a + banned native author sees their posts rejected and the community's own + moderation surface shows no ban. Deliberate scope cut for Scope A; + recorded so it is a decision rather than an omission. + ## Moderation state (task 17c-2) - **Bridge-side comment removal is WRITE-ONLY.** Nothing reads diff --git a/internal/accept/engine.go b/internal/accept/engine.go index 99e91d1..c569aac 100644 --- a/internal/accept/engine.go +++ b/internal/accept/engine.go @@ -100,6 +100,14 @@ const ( // that records "removed / moderator-removed" answers 200 accepted/enqueued. var ErrModeratorRemovalStands = stderrors.New("accept: a moderator removal stands") +// ErrAuthorBanned reports that the acceptance transaction found a ban the +// admission gate had not seen — the community banned this author between the +// two. Like the removal sentinel it is a DECISION rather than a failure: the +// transaction rolls back (so no acceptance, no outbound row, no delivery), the +// ledger records the rejection, and the live path treats the event as handled. +// Retrying would re-decide against a ban that is not going to move. +var ErrAuthorBanned = stderrors.New("accept: the author is banned from this community") + // errRemovalChanged reports that the removal an edit was deciding about is no // longer the record it inspected — it vanished, or a different one replaced it. // The decision is re-run against whatever stands now: an edit may only reverse @@ -256,6 +264,18 @@ func NewEngine(opts Options) (*Engine, error) { bans = fromCommunities } } + if bans == nil { + // REQUIRED, like the echo classifier and for the same reason: this is a + // gate, and a gate that is absent does not fail — it ADMITS. An engine + // built with a Communities that is not ban-capable (a fake, a decorator, + // a future backend) would accept every banned author's post into the + // community that excluded them, silently, with no counter and no log, + // and the deployment that did it would look identical to a correct one. + // The write side already refuses loudly when it cannot record a ban; the + // read side must not be the lenient half of the same feature. + return nil, errors.NewValidationError("bans", + "must not be nil: pass a store.CommunityBans, or a Communities that provides one") + } return &Engine{ repos: opts.Repos, enqueuer: opts.Enqueuer, @@ -337,6 +357,33 @@ func (e *Engine) AdmitPost(ctx context.Context, did string, commit *consume.Comm // federated once, so leaving it alone would strand it live on Lemmy); one // that was never accepted is simply a recorded rejection. priorAccepted := priorBound && !prior.IsTombstoned() + if priorAccepted && code == DecisionAuthorBanned { + // EXCEPT for a ban, which is not a judgement of this post. Removing + // here would do two things the ban itself deliberately did not: + // strip content Lemmy KEPT (a ban without removeData leaves it + // standing on their side, so we would be hiding a post they still + // show), and ENQUEUE a Delete{Page} at the community that banned the + // author — the outbound echo every other moderation path exists to + // avoid. It would also record RemovalCodeAdmissionRevoked, whose + // meaning is "a corrective edit may restore this", against a decision + // documented to survive an unban: one cause, two removals, opposite + // reversals, depending on which door the author knocked on. + // + // The ban already stopped everything new and cancelled everything + // queued. The edit is simply refused, and the ledger says why. + e.logger.Info("refusing a banned author's edit; the standing acceptance is left alone", + slog.String("did", did), slog.String("post", postURI), + slog.String("community", communityDID)) + return e.admissions.Record(ctx, Admission{ + AuthorDID: did, + CommunityDID: communityDID, + PostURI: postURI, + Status: StatusRejected, + DecisionCode: code, + EvaluatedCID: commit.CID, + EvaluatedSnapshot: e.evaluatedSnapshot(commit), + }) + } if priorAccepted { return e.removeAccepted(ctx, did, communityDID, postURI, commit, prior, code) } @@ -357,6 +404,24 @@ func (e *Engine) AdmitPost(ctx context.Context, did string, commit *consume.Comm } if err := e.accept(ctx, did, communityDID, postURI, commit); err != nil { + if stderrors.Is(err, ErrAuthorBanned) { + // The ban landed between the gate and the commit. The transaction + // rolled back, so nothing of this post exists outward; all that is + // owed is the ledger row an operator reads when the moderators ask + // why a banned author's post appeared — which it now will not. + e.logger.Info("a ban landed mid-admission; the post was not accepted", + slog.String("did", did), slog.String("post", postURI), + slog.String("community", communityDID)) + return e.admissions.Record(ctx, Admission{ + AuthorDID: did, + CommunityDID: communityDID, + PostURI: postURI, + Status: StatusRejected, + DecisionCode: DecisionAuthorBanned, + EvaluatedCID: commit.CID, + EvaluatedSnapshot: e.evaluatedSnapshot(commit), + }) + } if stderrors.Is(err, ErrModeratorRemovalStands) { // Decided and recorded inside accept(): the community removed this // post, so the edit does not re-enter it. Nothing is owed on the @@ -425,17 +490,14 @@ func (e *Engine) decide(ctx context.Context, did string, commit *consume.CommitE // decision about its own space, and admitting the post would sign that // community's name to content from someone it has excluded. // - // A nil store is "this deployment records no bans", not "nobody is banned": - // production wires it and NewEngine defaults it off the communities store, - // so the nil is only reachable from a caller that passes neither. - if e.bans != nil { - banned, err := e.bans.Standing(ctx, communityDID, did) - if err != nil { - return "", false, fmt.Errorf("accept: read ban on %s in %s: %w", did, communityDID, err) - } - if banned { - return DecisionAuthorBanned, false, nil - } + // NewEngine refuses to build without a ban store, so this is never a + // conditional check that quietly does not run. + banned, err := e.bans.Standing(ctx, communityDID, did) + if err != nil { + return "", false, fmt.Errorf("accept: read ban on %s in %s: %w", did, communityDID, err) + } + if banned { + return DecisionAuthorBanned, false, nil } // 5. Opt-out (decision 11): content pushed outward is exactly what an @@ -547,6 +609,26 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, // which rolls the acceptance back too (ApplyOpsTx side-effect atomicity), and // AdmitPost propagates it so the event retries. sideEffect := func(sctx context.Context, tx *sql.Tx, _ *repo.CommitResult) error { + // THE BAN IS RE-ASKED HERE, inside the transaction that acts on the + // answer. decide() read it before this transaction existed, and a ban is + // not just a row: the transaction that writes one also CANCELS every + // pending delivery the author has for this community. So a post that + // passed the gate and then commits behind the ban lands an acceptance — + // the community's own endorsement, in its own repo — plus a fresh + // delivery the cancellation could never have caught, aimed at the + // instance that just banned its author. + // + // It narrows the window rather than closing it: under READ COMMITTED a + // ban committing after this read and before this commit is still + // possible. What remains is microseconds wide and self-correcting on the + // next edit, where the old shape was seconds wide and permanent. + banned, err := e.bans.StandingTx(sctx, tx, communityDID, did) + if err != nil { + return fmt.Errorf("accept: re-read ban on %s in %s: %w", did, communityDID, err) + } + if banned { + return fmt.Errorf("%w: %s in %s", ErrAuthorBanned, postURI, communityDID) + } stored, err := e.objects.UpsertTx(sctx, tx, store.OutboundObject{ ATURI: postURI, APObjectID: apObjectID, diff --git a/internal/accept/outer_acceptance_test.go b/internal/accept/outer_acceptance_test.go index 37dd331..956daac 100644 --- a/internal/accept/outer_acceptance_test.go +++ b/internal/accept/outer_acceptance_test.go @@ -220,7 +220,14 @@ func acceptanceDB(t *testing.T) *sql.DB { testutil.Truncate(t, database, "ap_actors", "ap_objects", "communities", "repo_state", "blocks", "firehose_events", "outbound_activities", "outbound_deliveries", "outbound_objects", "outbound_votes", - "federation_prefs", "jetstream_record_revs", "jetstream_dead_letters", "admissions") + "federation_prefs", "jetstream_record_revs", "jetstream_dead_letters", "admissions", + // community_bans (migration 027). THIRD package to need this line, and + // the failure never names it: a ban left by one test makes decide() + // reject the NEXT test's author, which surfaces as "no acceptance + // record" and "no actor was minted" — the engine looking broken rather + // than the fixture being dirty. Any table a gate READS belongs here the + // moment a test WRITES it. + "community_bans") return database } diff --git a/internal/accept/terminality_race_test.go b/internal/accept/terminality_race_test.go index 27bca3e..6b6e4a3 100644 --- a/internal/accept/terminality_race_test.go +++ b/internal/accept/terminality_race_test.go @@ -2,6 +2,7 @@ package accept import ( "context" + "database/sql" "sync" "testing" "time" @@ -11,6 +12,7 @@ import ( "tidepool/internal/acceptrec" "tidepool/internal/repo" + "tidepool/internal/store" ) // TERMINALITY IS DECIDED ACROSS TWO OPERATIONS, AND THE GAP IS THE BUG. @@ -51,6 +53,13 @@ type racingRepos struct { // read and the restore. beforeRemovalDelete func() deleteFired bool + + // beforeAcceptanceWrite fires once, immediately before the commit that + // WRITES an acceptance — the moment after admission has decided and before + // anything of that decision is durable. It is where a ban lands in the + // window between the gate reading it and the acceptance committing. + beforeAcceptanceWrite func() + acceptanceFired bool } func (r *racingRepos) GetRecord(ctx context.Context, did, collection, rkey string) (map[string]any, string, error) { @@ -95,16 +104,34 @@ func (r *racingRepos) ApplyOpsTx(ctx context.Context, did string, ops []repo.Rec deletesRemoval = true } } + // The acceptance WRITE is the other side of the same window: any op that + // puts an acceptance record (never a delete of one) is the commit admission + // has already decided on. + writesAcceptance := false + for _, op := range ops { + if op.Action != repo.OpActionDelete && op.Collection == acceptrec.CollectionAcceptance { + writesAcceptance = true + } + } + r.mu.Lock() fire := deletesRemoval && r.beforeRemovalDelete != nil && !r.deleteFired if fire { r.deleteFired = true } hook := r.beforeRemovalDelete + fireAccept := writesAcceptance && r.beforeAcceptanceWrite != nil && !r.acceptanceFired + if fireAccept { + r.acceptanceFired = true + } + acceptHook := r.beforeAcceptanceWrite r.mu.Unlock() if fire { hook() } + if fireAccept { + acceptHook() + } return r.RepoManager.ApplyOpsTx(ctx, did, ops, sideEffect) } @@ -231,3 +258,104 @@ func TestVanishedRemovalIsNotRecordedAsTerminal(t *testing.T) { "moderators reinstated stays invisible until its author edits again") assert.Equal(t, acPostCID2, cid, "pinning the version the redrive evaluated") } + +// TASK 17c-3 REVIEW, P1-a — THE BAN GATE IS READ OUTSIDE THE TRANSACTION IT +// PROTECTS. +// +// Admission asks Standing() before opening the transaction that writes the +// acceptance and enqueues the delivery. Between those two moments a ban can +// land — and a ban is not just a row: the same transaction cancels every pending +// delivery the author has for that community. So the interleaving is +// +// post: Standing() → no ban +// ban: INSERT the row + CANCEL every pending delivery (commits) +// post: write the acceptance + enqueue a NEW delivery (commits) +// +// and the post lands AFTER the cancellation that existed to stop exactly it. +// The result is banned content accepted into the community's repo — visible in +// Coves under that community's name — and a fresh pending delivery carrying it +// to the instance that just banned its author, where it will be rejected, +// retried and poisoned. +// +// This is 17c-1's TERM-1 one verb over: a decision read outside the transaction +// that acts on it. The window is small and entirely ordinary — a moderator bans +// someone mid-thread while their client is uploading the next post. +func TestAPostCannotSlipPastABanThatLandsMidAdmission(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + repos := newRepos(t, conn) + seedBridgedCommunity(t, conn) + + racing := &racingRepos{RepoManager: repos} + enqueuer := realEnqueuer(t, conn) + engine := wireEngine(t, conn, racing, enqueuer) + dispatcher := wireDispatcher(t, conn, engine, enqueuer) + + // An earlier accepted post, so the ban has queued work to cancel. Without it + // the ban is only a row, and the race would be about visibility alone — the + // point here is that the post slips past a cancellation that already ran. + admittedCreate(t, dispatcher) + require.NotZero(t, pendingDeliveries(t, conn, acAuthorDID), + "precondition: the author has queued work for this community") + + bans := store.NewCommunityBans(conn) + racing.beforeAcceptanceWrite = func() { + cancelled, err := bans.Ban(ctx, store.CommunityBan{ + CommunityDID: acCommunityDID, + SubjectDID: acAuthorDID, + CommunityAPID: acCommunityAPID, + Reason: "banned while their next post was in flight", + }) + require.NoError(t, err, "the community's ban lands mid-admission") + require.NotZero(t, cancelled, + "and it cancels the work already queued — the state the racing post must not be "+ + "admitted behind") + } + + const racingRKey = "3lzpostban001" + racingURI := "at://" + acAuthorDID + "/social.coves.community.postv2/" + racingRKey + createsBefore := activityKindCount(t, conn, "Create") + + admitErr := dispatcher.HandleEvent(ctx, + postEvent("create", racingRKey, "3lzpostrev900", acPostCID, acPostTimeUS+10, pv2Record())) + require.True(t, racing.acceptanceFired, + "the fixture must have raced the acceptance commit, or this test is asserting nothing") + // Logged, not asserted: whether the engine reports success is not the + // property — an admission that returns nil while the ban stands is exactly + // the shape nothing upstream notices. + t.Logf("the racing admission reported: %v", admitErr) + + _, accepted := acceptanceSubjectCID(t, repos, acCommunityDID, racingURI) + assert.False(t, accepted, + "no acceptance may be written for a post the community has banned its author over: "+ + "the acceptance IS the community's endorsement, published in its own repo under "+ + "its own key, and it appears in Coves moments after the moderators excluded them") + + assert.Zero(t, pendingDeliveries(t, conn, acAuthorDID), + "and NOTHING may be left pending: the ban cancelled this author's queue, so a "+ + "delivery enqueued after that is one the cancellation could never have caught — "+ + "it goes to the instance that just banned them, is rejected, retries, and poisons") + + assert.Equal(t, createsBefore, activityKindCount(t, conn, "Create"), + "nor may a new Create activity exist for it at all") + + status, code := admissionOf(t, conn, acCommunityDID, racingURI) + assert.NotEqual(t, StatusAccepted, status, + "and the ledger must not record it as accepted: that row is what an operator reads "+ + "when the moderators ask why a banned author's post is in their community") + if status != StatusAccepted { + assert.NotEmpty(t, code, "with a machine-readable why") + } +} + +// pendingDeliveries counts an actor's queued outbound work. +func pendingDeliveries(t *testing.T, conn *sql.DB, actorDID string) int { + t.Helper() + var n int + require.NoError(t, conn.QueryRowContext(context.Background(), ` + SELECT COUNT(*) + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE a.actor_did = $1 AND d.state = 'pending'`, actorDID).Scan(&n)) + return n +} diff --git a/internal/ap/vocab.go b/internal/ap/vocab.go index b531c4c..dc73047 100644 --- a/internal/ap/vocab.go +++ b/internal/ap/vocab.go @@ -146,12 +146,18 @@ type Object struct { // of language objects on Group actors; Languages accepts both. Language Languages `json:"language,omitempty"` - // Expires is a Block's ban expiry. It is LOAD-BEARING and not decoration: - // Lemmy sends NO Undo when a temporary ban lapses — the ban simply stops - // applying on their side — so an implementation that drops this column turns - // every timed ban into a permanent one with no activity that can ever clear - // it. + // Expires and EndTime are the two spellings of a Block's ban expiry. Lemmy + // 0.19 sends `expires`; newer versions send AS2's `endTime` for the same + // fact. Both are kept as they arrived rather than merged at parse time — + // the wire said what it said — and BanExpiry answers which one applies. + // + // The field is LOAD-BEARING and not decoration: Lemmy sends NO Undo when a + // temporary ban lapses (it simply stops applying on their side), so an + // implementation that misses it turns every timed ban into a permanent one + // with no activity that could ever clear it. Reading only one spelling is + // exactly that bug, arriving on a version upgrade. Expires *Time `json:"expires,omitempty"` + EndTime *Time `json:"endTime,omitempty"` // Lemmy extensions. // RemoveData is Block's purge flag: the moderator also removed that author's @@ -302,6 +308,23 @@ func (o Object) isIDOnly() bool { return reflect.DeepEqual(clone, Object{}) } +// BanExpiry is the expiry a Block carries, under whichever spelling the sender +// used (`expires`, or AS2's `endTime` from newer Lemmy). nil means the activity +// carried NO expiry — a permanent ban — which is a different fact from an +// expiry that was present and could not be parsed: that one comes back non-nil +// with Valid false, and callers MUST tell the two apart. Collapsing them maps +// "we could not read how long" onto "forever", on the one field where forever +// is unrecoverable. +func (o *Object) BanExpiry() *Time { + if o == nil { + return nil + } + if o.Expires != nil { + return o.Expires + } + return o.EndTime +} + // IsActor reports whether the object's type is an AP actor type. func (o *Object) IsActor() bool { switch o.Type { diff --git a/internal/consume/comments.go b/internal/consume/comments.go index 78dd19c..bb06d20 100644 --- a/internal/consume/comments.go +++ b/internal/consume/comments.go @@ -102,6 +102,10 @@ func (d *Dispatcher) applyCommentWrite(ctx context.Context, tx *sql.Tx, did stri return err } + if err := d.refuseBannedAuthor(ctx, did, thread.CommunityDID, "comment "+atURI); err != nil { + return err + } + if err := d.ensureActor(ctx, did); err != nil { return err } @@ -131,6 +135,45 @@ func (d *Dispatcher) applyCommentWrite(ctx context.Context, tx *sql.Tx, did stri return d.enqueueComment(ctx, tx, did, commit.Operation, stored, thread.ParentATURI, thread.ParentAPID) } +// refuseBannedAuthor refuses anything a BANNED author writes into the community +// that banned them (task 17c-3 review). +// +// The admission gate covers POSTS only — it lives in the acceptance engine, +// which never sees a comment or a vote — so without this a banned author's +// replies and votes keep enqueueing on that community's ordering key. Lemmy +// refuses a banned actor's activity, so each one fails, retries and poisons with +// the cause three tables away: the same harm refuseInLockedThread exists to +// prevent, over a cause that is even better known here, because WE recorded it. +// +// The ban's own arrival already swept this author's queued comments and votes — +// its cancellation joins on actor_did and ordering_key, not on activity kind — +// so only NEW writes escaped, which is exactly what this closes. +// +// PERMANENT, like the lock refusal and for the same reasons: the event +// dead-letters with a reason an operator can read, the cursor advances rather +// than blocking every other native user behind one banned account, and only a +// moderator lifting the ban (or its expiry) changes the answer. +// +// The `what` argument is the short label the DLQ shows ("comment at://…", "vote +// at://…"): both paths share this refusal, and last_error is the only place the +// queue says which kind of write was refused. +func (d *Dispatcher) refuseBannedAuthor(ctx context.Context, did, communityDID, what string) error { + if communityDID == "" { + // Nothing resolved a community, so there is no ban to ask about: the + // caller has already decided this write goes nowhere. + return nil + } + banned, err := d.bans.Standing(ctx, communityDID, did) + if err != nil { + return fmt.Errorf("read ban on %s in %s: %w", did, communityDID, err) + } + if !banned { + return nil + } + return fmt.Errorf("%w: author-banned: %s is by an author %s has banned", + ErrPermanentEvent, what, communityDID) +} + // refuseInLockedThread refuses a comment in a thread a community has LOCKED // (task 17c-2). It is the read half of the lock: recording one and then // federating a reply under it is worse than not recording it at all, because diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index 2b85f55..8def8a1 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -195,6 +195,11 @@ type Options struct { Votes store.OutboundVotes Communities store.Communities ObjectMappings store.APObjects + // Bans reads whether the community a comment or vote is bound for has banned + // its author (task 17c-3 review). The acceptance engine gates POSTS; nothing + // else does, and a banned author's replies and votes enqueue to a community + // that refuses them until each delivery poisons. + Bans store.CommunityBans // Moderation is the bridge-owned moderation state the comment path reads to // refuse a reply in a locked thread. It is a SEPARATE store from // ObjectMappings on purpose: this consumer only ever reads it, while the @@ -229,6 +234,7 @@ type Dispatcher struct { votes store.OutboundVotes communities store.Communities moderation store.ObjectModeration + bans store.CommunityBans records materialize.RecordGetter hosted *hostedRepos gate *RevGate @@ -299,6 +305,7 @@ func NewDispatcher(opts Options) (*Dispatcher, error) { votes: orDefault[store.OutboundVotes](opts.Votes, store.NewOutboundVotes(opts.DB)), communities: orDefault[store.Communities](opts.Communities, store.NewCommunities(opts.DB)), moderation: orDefault[store.ObjectModeration](opts.Moderation, store.NewObjectModeration(opts.DB)), + bans: orDefault[store.CommunityBans](opts.Bans, store.NewCommunityBans(opts.DB)), records: opts.Records, hosted: newHostedRepos(opts.DB), gate: NewRevGate(opts.DB), diff --git a/internal/consume/dispatch_test.go b/internal/consume/dispatch_test.go index 683a8ad..444ee1f 100644 --- a/internal/consume/dispatch_test.go +++ b/internal/consume/dispatch_test.go @@ -31,7 +31,13 @@ func dispatchTestDB(t *testing.T) *sql.DB { testutil.Truncate(t, database, "ap_actors", "ap_objects", "communities", "repo_state", "outbound_objects", "outbound_votes", "federation_prefs", - "consumer_cursors", "jetstream_record_revs", "jetstream_dead_letters") + "consumer_cursors", "jetstream_record_revs", "jetstream_dead_letters", + // The gates the comment and vote paths now read (17c-2 locks, 17c-3 + // bans). FIFTH package to need this: any table a GATE READS belongs + // here the moment any test WRITES it, or a leftover row refuses the + // next test's author and the failure names neither the ban nor the + // test that left one. + "object_moderation", "community_bans") return database } diff --git a/internal/consume/outer_acceptance_test.go b/internal/consume/outer_acceptance_test.go index b06cd64..71a4e33 100644 --- a/internal/consume/outer_acceptance_test.go +++ b/internal/consume/outer_acceptance_test.go @@ -227,7 +227,13 @@ func consumeTestDB(t *testing.T) *sql.DB { testutil.Truncate(t, database, "ap_actors", "ap_objects", "communities", "repo_state", "outbound_objects", "outbound_votes", "federation_prefs", - "consumer_cursors", "jetstream_record_revs", "jetstream_dead_letters") + "consumer_cursors", "jetstream_record_revs", "jetstream_dead_letters", + // The gates the comment and vote paths now read (17c-2 locks, 17c-3 + // bans). FIFTH package to need this: any table a GATE READS belongs + // here the moment any test WRITES it, or a leftover row refuses the + // next test's author and the failure names neither the ban nor the + // test that left one. + "object_moderation", "community_bans") return database } diff --git a/internal/consume/votes.go b/internal/consume/votes.go index a5a7d0f..7d6ce22 100644 --- a/internal/consume/votes.go +++ b/internal/consume/votes.go @@ -88,6 +88,10 @@ func (d *Dispatcher) applyVoteWrite(ctx context.Context, tx *sql.Tx, did string, return nil } + if err := d.refuseBannedAuthor(ctx, did, subject.CommunityDID, "vote "+commitRecordURI(did, commit)); err != nil { + return err + } + if err := d.ensureActor(ctx, did); err != nil { return err } diff --git a/internal/db/migrations/027_community_bans.sql b/internal/db/migrations/027_community_bans.sql index 316669d..1843bfe 100644 --- a/internal/db/migrations/027_community_bans.sql +++ b/internal/db/migrations/027_community_bans.sql @@ -10,13 +10,19 @@ CREATE TABLE community_bans ( community_did TEXT NOT NULL, -- the bridged community that issued the ban subject_did TEXT NOT NULL, -- the native author it excludes - -- community_ap_id is DENORMALIZED, and it is not tidiness: the two readers - -- hold different handles for the same community. The admission gate has the - -- community DID (it is deciding about a repo it writes into); the delivery - -- queue has only OrderingKey, which IS the AP group id. A join the delivery - -- side cannot make is a scope it cannot apply — and it can never drift, - -- because the DID↔group-id mapping is immutable 1:1 (store.Communities - -- rejects an upsert that moves either). + -- community_ap_id is DENORMALIZED, and the honest status is: ONE reader + -- today, and it is this table's own write path. Ban() cancels the author's + -- pending deliveries, which are keyed by ordering_key — the AP group id, not + -- the DID — so the column is what makes the cancellation expressible in the + -- same statement as the row. + -- + -- It is kept for a reader that does not exist yet, deliberately: a + -- delivery-side ban recheck (claim time, where the worker holds ActorDID and + -- OrderingKey and nothing else) is the last line of defence against work + -- queued in the window before a ban lands. That reader would need an index on + -- (subject_did, community_ap_id); none exists, because nothing reads it that + -- way yet. It cannot drift meanwhile — the DID↔group-id mapping is immutable + -- 1:1 (store.Communities rejects an upsert that moves either). community_ap_id TEXT NOT NULL, banned_at TIMESTAMPTZ NOT NULL DEFAULT now(), -- expires_at is MANDATORY to honour, not optional to store. Lemmy's diff --git a/internal/ingest/consent.go b/internal/ingest/consent.go index 9acb220..d511daa 100644 --- a/internal/ingest/consent.go +++ b/internal/ingest/consent.go @@ -325,6 +325,15 @@ func (h *Handler) handleUndo(ctx context.Context, undo *ap.Object, signer string // only one Lemmy will ever send about it. return h.handleLock(ctx, inner, announcer, false) case ap.TypeBlock: + if announcer == nil { + // Lemmy sends BOTH halves of a ban to the banned user's inbox + // directly. Only the announced copies can be authorized, so both + // direct copies are ignored — and both are COUNTED, on the same + // counter: a number that moves for direct bans but not direct unbans + // makes a broken announce path look like a community that bans and + // never forgives. + return h.ignoreDirectBlock(inner, signer) + } // The unban, with the Block carried INLINE. It lifts the exclusion and // nothing else: content removed under removeData stays removed, because // Lemmy models restoration as a separate restore_data flag. diff --git a/internal/ingest/handler.go b/internal/ingest/handler.go index 5f606c0..6709e90 100644 --- a/internal/ingest/handler.go +++ b/internal/ingest/handler.go @@ -387,10 +387,10 @@ func (h *Handler) handleAnnounce(ctx context.Context, announce *ap.Object, signe // moderator's Person and cannot satisfy decision 18 by construction. return h.handleBlock(ctx, inner, community, true) default: - // Add, Remove, Block, ... — moderation activities the bridge does not + // Add, Remove, Flag, ... — moderation activities the bridge does not // translate yet. Remove in particular is NOT content removal in Lemmy // (it is un-pin / demote-moderator, dispatched by `target`), so it must - // never be folded in beside Lock on the assumption that it is. + // never be folded in beside Lock or Block on the assumption that it is. return skip(announce.ID, "unsupported announced activity type "+inner.Type) } } diff --git a/internal/ingest/moderation.go b/internal/ingest/moderation.go index 190023d..89a8620 100644 --- a/internal/ingest/moderation.go +++ b/internal/ingest/moderation.go @@ -4,6 +4,7 @@ import ( "context" "expvar" "fmt" + "time" "tidepool/internal/ap" "tidepool/internal/echo" @@ -45,8 +46,21 @@ var ( // of our personas. A Lemmy user banned from a Lemmy community is entirely // their instance's business; we hold no state that could apply it. BlockForeignSubject = expvar.NewInt("tidepool_block_foreign_subject") + // BlockLapsedIgnored counts announced Blocks whose expiry had already passed + // when they arrived — a description of a ban that is over, applied to + // nothing. + BlockLapsedIgnored = expvar.NewInt("tidepool_block_lapsed_ignored") + // BlockExpiryUnreadable counts Blocks refused because their expiry could not + // be parsed. It is the counter that says "a peer is sending us a duration we + // do not understand" — which, unlike most parse failures, would otherwise + // have become a permanent ban. + BlockExpiryUnreadable = expvar.NewInt("tidepool_block_expiry_unreadable") ) +// timeNow is the clock the ban path weighs an expiry against. A package +// variable so a test can hold time still; production never replaces it. +var timeNow = time.Now + // ignoreDirectBlock is the DECIDED non-action at the other door. // // Lemmy sends a ban twice: announced through the community, and delivered @@ -177,17 +191,56 @@ func (h *Handler) applyBan(ctx context.Context, block *ap.Object, announcer *sto // event processed would leave the community believing we honoured it. return fmt.Errorf("ingest: no community-ban store is wired, so this ban cannot be recorded") } + // THE EXPIRY IS DECIDED BEFORE ANY CONSEQUENCE, because every consequence + // below is irreversible: cancelled deliveries are never re-queued, and + // removeData's removals are terminal by design (no Undo{Block} restores + // content). A ban's duration therefore has to be settled while doing nothing + // is still an option. + expiry := block.BanExpiry() + switch { + case expiry == nil: + // No expiry: a permanent ban, which is the common case. + case !expiry.Valid: + // PRESENT BUT UNREADABLE. The parser keeps this apart from absent + // precisely so it can be refused: treating it as "no expiry" records a + // permanent exclusion the moderator did not ask for, and nothing would + // ever correct it — Lemmy sends no activity when a ban lapses, so there + // is no later message whose arrival could say "that should have ended". + // From every side it would read as an ordinary permanent ban. + // + // A validation error POISONS rather than retries: the bytes will not + // re-parse, and the honest outcome is a visible failure saying we did not + // apply this ban, not a queue that re-reads the same string forever. + BlockExpiryUnreadable.Add(1) + return errors.NewValidationError("expires", + "block for "+subjectDID+" carries an expiry that cannot be parsed; refusing to "+ + "store it as a permanent ban") + } + // ALREADY OVER when it arrived — delayed in a queue, redelivered after an + // outage, replayed from a backfill. The ROW is still written below: it is a + // faithful account of what the moderator sent, it makes a redelivery + // idempotent, and Standing() reads the expiry so it excludes nobody. + // + // What a lapsed ban must NOT do is ACT. Every consequence here is one no + // later activity can undo — a cancelled delivery is never re-queued, and a + // removeData removal is terminal by design — so applying them over an + // exclusion that has already ended is unrecoverable damage done on behalf of + // a decision that expired. The cancellation is gated inside Ban() (one place, + // beside the statement); the purge is gated here. + lapsed := expiry != nil && expiry.Valid && !expiry.After(timeNow()) + ban := store.CommunityBan{ CommunityDID: announcer.DID, SubjectDID: subjectDID, CommunityAPID: announcer.APGroupID, - RemoveData: block.RemoveData != nil && *block.RemoveData, - } - // The expiry is carried through EXACTLY as sent. Lemmy sends no activity when - // a timed ban lapses — it simply stops applying there — so dropping this - // makes a three-day ban permanent with nothing that could ever clear it. - if block.Expires.OK() { - expires := block.Expires.Time + // The moderator's own words, kept because the SAME action already writes + // them into any removal record it produces: a blank here beside a quoted + // reason there tells an operator no reason was given. + Reason: block.Summary, + RemoveData: block.RemoveData != nil && *block.RemoveData, + } + if expiry != nil { + expires := expiry.Time ban.ExpiresAt = &expires } @@ -201,8 +254,15 @@ func (h *Handler) applyBan(ctx context.Context, block *ap.Object, announcer *sto h.logger.Info("community banned a native author", "community", announcer.APGroupID, "subject_did", subjectDID, "expires", ban.ExpiresAt, "remove_data", ban.RemoveData, - "cancelled_deliveries", cancelled, "activity", block.ID) + "lapsed", lapsed, "cancelled_deliveries", cancelled, "activity", block.ID) + if lapsed { + BlockLapsedIgnored.Add(1) + return skip(block.ID, + "announced Block expired before it arrived: the ban is recorded as sent, but it "+ + "is not in force — nothing was cancelled and nothing was removed, because both "+ + "are irreversible and this exclusion is already over") + } if !ban.RemoveData { return nil } diff --git a/internal/ingest/moderation_ban_test.go b/internal/ingest/moderation_ban_test.go index d434387..6d2821f 100644 --- a/internal/ingest/moderation_ban_test.go +++ b/internal/ingest/moderation_ban_test.go @@ -3,7 +3,10 @@ package ingest import ( "context" "database/sql" + "encoding/json" + stderrors "errors" "expvar" + "fmt" "net/http" "testing" "time" @@ -13,6 +16,7 @@ import ( "tidepool/internal/accept" "tidepool/internal/ap" + "tidepool/internal/consume" "tidepool/internal/errors" "tidepool/internal/materialize" ) @@ -252,16 +256,18 @@ func blockActivity(group *remoteActor, activityID, subjectDID, target string, ex type communityBan struct { communityAPID string expires sql.NullTime + reason string + removeData bool } func banFor(t *testing.T, db *sql.DB, communityDID, subjectDID string) (communityBan, bool) { t.Helper() var ban communityBan err := db.QueryRowContext(context.Background(), ` - SELECT community_ap_id, expires_at + SELECT community_ap_id, expires_at, reason, remove_data FROM community_bans WHERE community_did = $1 AND subject_did = $2`, - communityDID, subjectDID).Scan(&ban.communityAPID, &ban.expires) + communityDID, subjectDID).Scan(&ban.communityAPID, &ban.expires, &ban.reason, &ban.removeData) if err == sql.ErrNoRows { return communityBan{}, false } @@ -288,6 +294,21 @@ func requireEveryDelivery(t *testing.T, db *sql.DB, actorDID, orderingKey, want, } } +// assertEveryDelivery is requireEveryDelivery without the abort. A test pinning +// SEVERAL independent consequences of one activity uses this so a broken first +// consequence does not hide the others — the vacuity guard stays fatal, because +// a test asserting over no rows has stopped measuring anything. +func assertEveryDelivery(t *testing.T, db *sql.DB, actorDID, orderingKey, want, why string) { + t.Helper() + states := deliveryStates(t, db, actorDID, orderingKey) + require.NotEmpty(t, states, + "%s — and there must BE deliveries to say that about: %s has no rows on %s at all", + why, actorDID, orderingKey) + for i, state := range states { + assert.Equal(t, want, state, "%s (delivery %d of %d)", why, i+1, len(states)) + } +} + // deliveryStates lists the states of one actor's deliveries on one community's // ordering key. // @@ -407,33 +428,81 @@ func TestADirectBlockIsIgnoredWhateverItClaims(t *testing.T) { "shadowban nobody issued and nobody can lift") } -// TestALapsedBanDoesNotRefuseAdmission is the half of `expires` that has no -// activity behind it. +// TestALapsedBanIsInert is the half of `expires` that has no activity behind it, +// and it covers ALL THREE things a lapsed Block must not do. // // Lemmy's BlockUser carries `expires` for a temporary ban, and when that ban -// lapses Lemmy sends NOTHING — no Undo, no second activity, nothing. The ban -// simply stops applying on their side. So an implementation that stores the ban -// and ignores the column turns every 3-day ban into a permanent one, and there -// is no message that will ever clear it: the author is excluded forever by a -// moderator who chose three days. -func TestALapsedBanDoesNotRefuseAdmission(t *testing.T) { +// lapses Lemmy sends NOTHING — no Undo, no second activity. The ban simply stops +// applying on their side. So a Block whose expiry has ALREADY PASSED when it +// reaches us — delayed in a queue, redelivered after an outage, replayed from a +// backfill — is a description of a ban that is already over. +// +// Reading it as live is destructive in two directions that no later activity can +// repair: it cancels queued work that was never banned, and with removeData it +// PERMANENTLY removes accepted posts. Both are unreachable from the admission +// gate, which is why an admission-only test of the lapsed case passes while both +// are broken — the expiry has to be weighed where the ban ACTS, not only where +// it is later read. +func TestALapsedBanIsInert(t *testing.T) { h := newHarness(t) ctx := context.Background() world := newModerationWorld(t, h) + // Queued work, and an accepted post: the two things a live ban destroys. + admitPost(t, world, mtAuthorDID, mbPostInARKey, world.communityADID, "3lzmbrev00050", 1_775_000_009_000_001) + requireEveryDelivery(t, h.db, mtAuthorDID, groupID, "pending", + "precondition: the author has queued work for A") + _, _, err := h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + require.NoError(t, err, "precondition: and an accepted post in A") + + // A ban that ended before it arrived — carrying removeData, so every + // destructive branch is on the table. h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, - map[string]any{"expires": "2020-01-01T00:00:00Z"}) + map[string]any{"expires": "2020-01-01T00:00:00Z", "removeData": true}) + + // (1) It cancels nothing. + assertEveryDelivery(t, h.db, mtAuthorDID, groupID, "pending", + "a LAPSED ban cancels nothing: these posts were queued by an author who is not "+ + "banned now and was not banned when the ban expired — cancelling them silently "+ + "unpublishes work on the strength of an exclusion that has already ended, and "+ + "nothing re-queues a cancelled delivery") + // (2) It removes nothing. + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionRemoval, world.digestRKey) + assert.True(t, errors.IsNotFound(err), + "and it removes NOTHING: removeData on a lapsed ban strips accepted posts for an "+ + "exclusion that is over, and a removal is terminal — no Undo{Block} restores "+ + "content, by design, so this is unrecoverable (err=%v)", err) + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + assert.NoError(t, err, "the acceptance stands (err=%v)", err) + + // (3) It refuses no admission — the half that was already covered, kept + // because the three only mean anything together. + // CONDITIONAL, and deliberately so: both outcomes are correct here. Not + // recording a lapsed ban at all is the stronger reading — a ban row is + // CURRENT STATE, not a log (Lift deletes rather than marks), so a row for an + // exclusion that is already over is state no reader wants and one more thing + // that has to remember the expiry filter. Recording it with its expiry is + // also sound, since every read is scoped by `expires_at > now()`. + // + // What this conditional CANNOT hide is the regression that matters: the only + // dangerous row shape is one with a NULL expiry, because that reads as + // permanent forever, and it fails inside the guard. An absent row cannot be + // mistaken for a standing ban by anything. if ban, found := banFor(t, h.db, world.communityADID, mtAuthorDID); found { require.True(t, ban.expires.Valid, - "a ban recorded from an activity carrying `expires` must carry the expiry with it: "+ - "dropping the column is what makes the lapse unrepresentable") - assert.True(t, ban.expires.Time.Before(timeNow()), + "a lapsed ban that IS recorded must carry its expiry: storing it unbounded turns "+ + "a ban that already ended into a permanent one, and no activity will ever "+ + "correct that — Lemmy sends nothing when a ban lapses") + assert.True(t, ban.expires.Time.Before(time.Now()), "and it must be the moment the moderator chose, in the past") } admitPost(t, world, mtAuthorDID, mbPostAfterBanRKey, world.communityADID, - "3lzmbrev00011", 1_775_000_003_000_001) + "3lzmbrev00051", 1_775_000_009_000_002) postURI := mbPostATURI(mtAuthorDID, mbPostAfterBanRKey) status, code := admissionFor(t, h.db, world.communityADID, postURI) assert.Equal(t, accept.StatusAccepted, status, @@ -441,11 +510,66 @@ func TestALapsedBanDoesNotRefuseAdmission(t *testing.T) { "expires_at > now()`, because no Undo is coming — the expiry IS the lift") assert.NotEqual(t, "author-banned", code, "and certainly not for being banned") - _, _, err := h.manager.GetRecord(ctx, + _, _, err = h.manager.GetRecord(ctx, world.communityADID, materialize.CollectionAcceptance, testDigestRKey(postURI)) assert.NoError(t, err, "with the acceptance to prove it (err=%v)", err) } +// TestABanWhoseExpiryCannotBeReadIsNotStoredAsPermanent is the fail-safe +// direction on the one field whose absence means FOREVER. +// +// `expires` absent and `expires` present-but-unparseable are different facts, +// and the wire parser keeps them apart (a nil *Time versus a Time with +// Valid=false) precisely so a caller can. Collapsing them maps "we could not +// read how long" onto "no expiry", which is the most consequential possible +// misreading: the moderator asked for a time limit, and the author is excluded +// forever instead. +// +// Nothing ever repairs it. Lemmy sends no activity when a timed ban lapses, so +// there is no message whose arrival could correct the record, and no operator +// has a reason to look — from every side this reads like an ordinary permanent +// ban that a moderator chose. +// +// The honest outcome is to refuse the activity rather than store a ban we cannot +// bound: we know what they meant and we could not read how long. +func TestABanWhoseExpiryCannotBeReadIsNotStoredAsPermanent(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + const malformed = "https://lemmy.world/activities/announce/block/mb-badexpiry" + h.announceBlock(world.groupA, malformed, mtAuthorDID, groupID, + map[string]any{"expires": "next tuesday"}) + + ban, found := banFor(t, h.db, world.communityADID, mtAuthorDID) + if found { + assert.True(t, ban.expires.Valid, + "a ban stored from an activity that CARRIED an expiry must carry one: storing it "+ + "with a NULL expiry records a permanent exclusion the moderator did not ask "+ + "for, and nothing will ever clear it") + } else { + assert.False(t, found, + "or it is not stored at all, which is the safer reading of an unbounded ban") + } + + event, err := h.events.GetEvent(ctx, malformed) + require.NoError(t, err) + assert.Nil(t, event.ProcessedAt, + "and the activity is NOT marked handled: an expiry we cannot read is a ban we cannot "+ + "bound, so the honest outcome is a failure an operator can see — retryable or "+ + "poisoned, either says 'we did not apply this'. Marking it processed tells the "+ + "community we honoured a ban whose duration we silently replaced with forever: %s", + event.Error) + + admitPost(t, world, mtAuthorDID, mbPostAfterBanRKey, world.communityADID, + "3lzmbrev00060", 1_775_000_010_000_001) + status, code := admissionFor(t, h.db, world.communityADID, mbPostATURI(mtAuthorDID, mbPostAfterBanRKey)) + assert.Equal(t, accept.StatusAccepted, status, + "and the author is not excluded on an unreadable instruction: the failure mode this "+ + "guards is a permanent shadowban created by a typo in someone else's software") + assert.NotEqual(t, "author-banned", code, "least of all as author-banned") +} + // TestAStandingTimedBanRefusesAdmission is the other half: the same activity // shape, an expiry that has NOT arrived, and a ban that bites. // @@ -463,7 +587,7 @@ func TestAStandingTimedBanRefusesAdmission(t *testing.T) { ban, found := banFor(t, h.db, world.communityADID, mtAuthorDID) require.True(t, found, "a temporary ban is still a ban and must be recorded") require.True(t, ban.expires.Valid, "carrying its expiry") - assert.True(t, ban.expires.Time.After(timeNow()), "which has not arrived") + assert.True(t, ban.expires.Time.After(time.Now()), "which has not arrived") admitPost(t, world, mtAuthorDID, mbPostAfterBanRKey, world.communityADID, "3lzmbrev00012", 1_775_000_004_000_001) @@ -529,6 +653,7 @@ func TestASiteScopedBlockIsSkippedWithItsOwnReason(t *testing.T) { // against another outcome, not against a substring. const siteActor = "https://lemmy.world/" const siteBlock = "https://lemmy.world/activities/announce/block/mb-site" + unscopedBefore := metricValue(mbUnscopedTargetMetric) h.announceBlock(world.groupA, siteBlock, mtAuthorDID, siteActor, nil) _, found := banFor(t, h.db, world.communityADID, mtAuthorDID) @@ -536,19 +661,43 @@ func TestASiteScopedBlockIsSkippedWithItsOwnReason(t *testing.T) { "a site ban is not a community ban: recording it against the announcing community "+ "would understate it — the user is excluded from every community on that instance, "+ "and we would enforce it in one") - - // The same activity, targeted at the community, is the shape that WORKS. - const communityBlock = "https://lemmy.world/activities/announce/block/mb-site-control" - h.announceBlock(world.groupA, communityBlock, mtCommenterDID, groupID, nil) - - assert.NotEqual(t, - skipReasonFor(t, h, communityBlock, communityBlock), - skipReasonFor(t, h, siteBlock, siteBlock), - "a target this scope does not model must be distinguishable from one it does: today "+ - "both land in the same 'unsupported activity type' default, so an operator whose "+ - "user is collecting 403s across a whole instance reads the same line as someone "+ - "whose ban was applied — and silently no-op'ing the site ban leaves that user "+ - "posting into the instance until their deliveries poison") + assert.Equal(t, unscopedBefore+1, metricValue(mbUnscopedTargetMetric), + "and it is COUNTED under its own class: this is the scope we chose not to model, and "+ + "the counter is what tells an operator how much of it is arriving — a user "+ + "collecting 403s across a whole instance is invisible otherwise") + + // THE CONTROL, ASSERTED: the same activity targeted at the announcer's own + // community WORKS. Without this the comparison below could be satisfied by a + // control that silently failed too. + const workingBlock = "https://lemmy.world/activities/announce/block/mb-site-control" + // The control's subject needs an AP persona to be bannable at all — the + // subject is resolved by ENTITY EXISTENCE, so an actor who has never + // federated anything is refused as a foreign subject, and the control would + // fail for a reason that has nothing to do with targets. + admitPost(t, world, mtCommenterDID, mbOtherPostRKey, world.communityADID, + "3lzmbrev00080", 1_775_000_012_000_001) + h.announceBlock(world.groupA, workingBlock, mtCommenterDID, groupID, nil) + _, worked := banFor(t, h.db, world.communityADID, mtCommenterDID) + require.True(t, worked, + "precondition: a community-targeted Block from the same announcer is recorded — the "+ + "shape under test differs from this one ONLY in its target") + + // And the discrimination is against another REFUSAL, not against a success: + // a successful Block logs no skip at all, so comparing with it would reduce + // to 'the site block produced some reason', which the generic + // unsupported-type default already satisfies. + const wrongCommunityBlock = "https://lemmy.world/activities/announce/block/mb-site-cross" + h.announceBlock(world.groupB, wrongCommunityBlock, mtAuthorDID, groupID, nil) + + siteReason := skipReasonFor(t, h, siteBlock, siteBlock) + crossReason := skipReasonFor(t, h, wrongCommunityBlock, wrongCommunityBlock) + require.NotEmpty(t, siteReason, "the site-targeted Block must be skipped with a reason") + require.NotEmpty(t, crossReason, "and so must the wrong-community one") + assert.NotEqual(t, crossReason, siteReason, + "a target this scope does not model must be distinguishable from a target that "+ + "belongs to somebody else: one means 'we do not implement instance bans' and the "+ + "other means 'that is not your community', and an operator who cannot tell them "+ + "apart cannot tell a missing feature from a rejected forgery") event, err := h.events.GetEvent(ctx, siteBlock) require.NoError(t, err) @@ -560,7 +709,15 @@ func TestASiteScopedBlockIsSkippedWithItsOwnReason(t *testing.T) { // NAME so the test survives the var being renamed or moved — the counter's // identity is its published name, which is what an operator's dashboard binds // to, not the Go symbol. -const mbDirectIgnoredMetric = "tidepool_block_direct_ignored" +const ( + mbDirectIgnoredMetric = "tidepool_block_direct_ignored" + // mbUnscopedTargetMetric counts announced Blocks whose target is not a + // community we follow — the instance-wide scope this bridge does not model. + mbUnscopedTargetMetric = "tidepool_block_unscoped_target" + // mbForeignSubjectMetric counts announced Blocks naming somebody who is not + // one of our personas. + mbForeignSubjectMetric = "tidepool_block_foreign_subject" +) func metricValue(name string) int64 { counter, _ := expvar.Get(name).(*expvar.Int) @@ -570,8 +727,6 @@ func metricValue(name string) int64 { return counter.Value() } -func timeNow() time.Time { return time.Now() } - // TestABanWithRemoveDataRemovesTheirPostsInThatCommunityOnly is the destructive // half, and its scope is the whole design. // @@ -693,6 +848,9 @@ func TestABanLeavesTerminalDeliveriesAlone(t *testing.T) { // about the ban rather than about delivery. terminal := forceDeliveryStates(t, h.db, mtAuthorDID, groupID, "delivered", "poisoned") require.Len(t, terminal, 2, "precondition: two terminal deliveries to ban across") + pendingBefore := pendingActivityIDs(t, h.db, mtAuthorDID, groupID) + require.NotEmpty(t, pendingBefore, + "precondition: and a PENDING one beside them, so the ban has something to cancel") h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, nil) @@ -704,6 +862,19 @@ func TestABanLeavesTerminalDeliveriesAlone(t *testing.T) { "and a POISONED row is untouched: sweeping it hides a delivery that failed for its "+ "own reason behind a ban that arrived afterwards, and redrive is how an operator "+ "gets it back") + + // THE BAN MUST HAVE DONE SOMETHING. Every assertion above is satisfied by a + // Block that was skipped entirely, which is the mis-attributed control in + // its purest form: a test that proves cancellation is correctly scoped by + // arranging for no cancellation to happen at all. + _, found := banFor(t, h.db, world.communityADID, mtAuthorDID) + require.True(t, found, "the ban was recorded — this test is about what it did NOT touch") + for _, id := range pendingBefore { + assert.Equal(t, "cancelled", deliveryState(t, h.db, id), + "...having cancelled the PENDING row beside them: if that one survived too, "+ + "nothing was cancelled and the two assertions above proved nothing about "+ + "scope — they would hold just as well for a Block that was skipped entirely") + } } // forceDeliveryStates stamps states onto an actor's pending deliveries for one @@ -738,6 +909,27 @@ func forceDeliveryStates(t *testing.T, db *sql.DB, actorDID, orderingKey string, return ids } +// pendingActivityIDs lists an actor's pending deliveries on one ordering key. +func pendingActivityIDs(t *testing.T, db *sql.DB, actorDID, orderingKey string) []string { + t.Helper() + rows, err := db.QueryContext(context.Background(), ` + SELECT d.activity_id + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE a.actor_did = $1 AND d.ordering_key = $2 AND d.state = 'pending' + ORDER BY d.seq`, actorDID, orderingKey) + require.NoError(t, err) + defer func() { require.NoError(t, rows.Close()) }() + var ids []string + for rows.Next() { + var id string + require.NoError(t, rows.Scan(&id)) + ids = append(ids, id) + } + require.NoError(t, rows.Err()) + return ids +} + func deliveryState(t *testing.T, db *sql.DB, activityID string) string { t.Helper() var state string @@ -800,3 +992,396 @@ func TestAnAnnouncedBlockLaunderedThroughAPersonIsRefused(t *testing.T) { "and admission is unaffected: a half-applied forged ban is a shadowban nobody issued "+ "and no Undo can lift, because no moderator ever made the decision to reverse") } + +// TASK 17c-3 REVIEW, P2 — WHAT THE BAN PATH RECORDS AND COUNTS. + +// TestADirectUndoBlockIsCountedLikeADirectBlock closes an asymmetry in the one +// number that reports on this door. +// +// Lemmy sends BOTH halves of a ban directly to the banned user's inbox: the +// Block, and later the Undo{Block}. Only the announced copies can be authorized, +// so both direct copies are ignored — but only the Block is counted. Direct +// unban traffic is therefore invisible, and the counter that is supposed to say +// "this is how much of the ban conversation arrives on the door we do not open" +// reports half of it. +// +// That matters precisely when it is read: if the announced path ever breaks, the +// operator compares direct traffic against recorded bans. A counter that moves +// for bans and not unbans makes a broken announce path look like a community +// that bans and never forgives. +func TestADirectUndoBlockIsCountedLikeADirectBlock(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + moderator := h.newRemoteActor(modActorID, person(modActorID, "moderator", nil)) + + before := metricValue(mbDirectIgnoredMetric) + + const directUndo = "https://lemmy.world/activities/undo/mb-direct-undo" + require.Equal(t, http.StatusAccepted, h.deliver(moderator, map[string]any{ + "id": directUndo, + "type": "Undo", + "actor": modActorID, + "audience": groupID, + "cc": []any{groupID}, + "object": blockActivity(world.groupA, directUndo+"/block", mtAuthorDID, groupID, nil), + })) + h.drain() + + assert.Equal(t, before+1, metricValue(mbDirectIgnoredMetric), + "a direct Undo{Block} is the same decided non-action as a direct Block and must be "+ + "counted the same way: Lemmy sends both to this door, and a counter that moves for "+ + "one and not the other reports the door as quieter than it is") + + event, err := h.events.GetEvent(ctx, directUndo) + require.NoError(t, err) + assert.NotNil(t, event.ProcessedAt, "decided once, not retried: %s", event.Error) + assert.Nil(t, event.FailedAt, "nor poisoned") +} + +// TestABanRecordsTheModeratorsReasonAndScope pins the two fields the row carries +// for people rather than for logic. +// +// `removeData` decides whether content goes, and the reason is what a moderator +// typed. The schema holds both and the repository writes both — but nothing +// fills them in, so every ban reads as reasonless. It matters because the SAME +// moderator action already writes the reason somewhere else: a removeData +// removal carries the summary into the removal record. Two audit surfaces +// describing one decision, one of them blank, is worse than neither having it — +// whoever reads the empty one concludes no reason was given. +func TestABanRecordsTheModeratorsReasonAndScope(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + + const reason = "repeated rule 3 violations after two warnings" + h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, + map[string]any{"summary": reason, "removeData": true}) + + ban, found := banFor(t, h.db, world.communityADID, mtAuthorDID) + require.True(t, found, "the ban is recorded") + assert.Equal(t, reason, ban.reason, + "with the moderator's own words: the column exists, the repository writes it, and the "+ + "same summary already reaches the removal records this ban produced — a blank here "+ + "tells an operator no reason was given while the removal beside it quotes one") + assert.True(t, ban.removeData, + "and with the scope it was issued under: whether content went is the difference "+ + "between an exclusion and a purge, and after the fact the row is the only place "+ + "that answers it") +} + +// TestABanCancelsWorkAWorkerHasAlreadyClaimed is the honest half of the +// in-flight problem. +// +// A worker claims a delivery in its own committed transaction and then POSTs. +// If a ban lands between those, nothing can un-send the request — the bytes are +// on the wire, and fencing only stops the settlement afterwards. So what is +// pinned here is what remains true and reachable: a CLAIMED row is still +// cancelled, and its claim is released rather than left to expire, so no further +// attempt is made for it. +// +// The unfixable remainder is a pre-send re-check, which this test deliberately +// does not pretend to cover. +func TestABanCancelsWorkAWorkerHasAlreadyClaimed(t *testing.T) { + h := newHarness(t) + world := newModerationWorld(t, h) + + admitPost(t, world, mtAuthorDID, mbPostInARKey, world.communityADID, "3lzmbrev00070", 1_775_000_011_000_001) + claimed := claimDeliveries(t, h.db, mtAuthorDID, groupID) + require.NotZero(t, claimed, "precondition: a worker holds a claim on this author's work") + + h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, nil) + + assertEveryDelivery(t, h.db, mtAuthorDID, groupID, "cancelled", + "a claimed row is cancelled like any other: the claim is a lease, not an exemption, "+ + "and leaving claimed work pending would let it be re-attempted for the whole lease "+ + "after the community banned its author") + assert.Zero(t, claimedDeliveries(t, h.db, mtAuthorDID, groupID), + "and the claim is RELEASED rather than left to expire: a cancelled row still holding "+ + "a lease is one a worker sweep has to reason about, and the point of cancelling is "+ + "that no further attempt is made") +} + +// claimDeliveries simulates a worker claim — a lease stamped in a transaction +// that has already committed, which is exactly the state a ban can arrive in the +// middle of. It returns how many rows it claimed. +func claimDeliveries(t *testing.T, db *sql.DB, actorDID, orderingKey string) int64 { + t.Helper() + result, err := db.ExecContext(context.Background(), ` + UPDATE outbound_deliveries d + SET claimed_until = now() + interval '1 minute' + FROM outbound_activities a + WHERE d.activity_id = a.activity_id + AND a.actor_did = $1 AND d.ordering_key = $2 AND d.state = 'pending'`, + actorDID, orderingKey) + require.NoError(t, err) + affected, err := result.RowsAffected() + require.NoError(t, err) + return affected +} + +func claimedDeliveries(t *testing.T, db *sql.DB, actorDID, orderingKey string) int { + t.Helper() + var n int + require.NoError(t, db.QueryRowContext(context.Background(), ` + SELECT COUNT(*) + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE a.actor_did = $1 AND d.ordering_key = $2 AND d.claimed_until IS NOT NULL`, + actorDID, orderingKey).Scan(&n)) + return n +} + +// TASK 17c-3 REVIEW, HIGH-1 — THE GATE COVERS POSTS ONLY. +// +// A ban excludes an AUTHOR from a community, and admission enforces it for +// postv2 alone. Comments gate on opt-out, thread, depth and the parent lock; +// votes gate on opt-out. So a banned author's replies and votes keep leaving for +// the community that banned them, Lemmy refuses each one, and the deliveries +// retry to poisoned with the cause three tables away — precisely the outcome the +// thread-lock refusal exists to prevent, and the same fix shape. +// +// The tell that this is a GAP and not a decision: the ban's own cancellation +// already sweeps this author's queued comments and votes. It joins actor_did and +// ordering_key, and never looks at what kind of activity a row carries. So the +// system already agrees that a banned author's comments and votes must not go to +// that community — it just enforces it on the traffic that happens to be queued +// when the ban lands, and not on anything written afterwards. + +// TestABannedAuthorsCommentIsRefusedInThatCommunityOnly is HIGH-1 for comments. +func TestABannedAuthorsCommentIsRefusedInThatCommunityOnly(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + // The same author, accepted in BOTH communities: the ban must reach exactly + // one of the two threads they can reply in. + admitPost(t, world, mtAuthorDID, mbPostInBRKey, world.communityBDID, "3lzmbrev00090", 1_775_000_013_000_001) + postInB := mbPostATURI(mtAuthorDID, mbPostInBRKey) + + h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, nil) + _, banned := banFor(t, h.db, world.communityADID, mtAuthorDID) + require.True(t, banned, "precondition: the author is banned from A and not from B") + + deliveriesBefore := rowCount(t, h.db, "outbound_deliveries") + + // --- Their comment into the community that banned them. + inA := mbCommentBy(mtAuthorDID, "3lzmbcomment01", mtPostATURI, mtPostCID, "3lzmbrev00091") + err := world.dispatcher.HandleEvent(ctx, inA.create(t)) + require.Error(t, err, + "a banned author's COMMENT must be refused: a ban excludes the author from the "+ + "community, not their postv2 records from admission — Lemmy rejects the comment "+ + "server-side, so federating it buys a failed delivery, a retry loop and a poisoned "+ + "row whose cause is a moderator decision nothing in the queue names") + assert.True(t, stderrors.Is(err, consume.ErrPermanentEvent), + "with the same permanence as a locked thread: retrying cannot change a moderator's "+ + "decision, and parking it holds every other native user behind it (err=%v)", err) + assert.Contains(t, err.Error(), "author-banned", + "and naming the cause in the one string the DLQ carries — the operator answering "+ + "'where did my comment go' has nothing else") + assert.Zero(t, outboundRowsFor(t, h.db, inA.atURI()), + "nothing is written for the refused comment") + assert.Equal(t, deliveriesBefore, rowCount(t, h.db, "outbound_deliveries"), + "and nothing is sent to a community that will reject it") + + // --- The same author, the same moment, the community they are NOT banned in. + inB := mbCommentBy(mtAuthorDID, "3lzmbcomment02", postInB, mtPostCID, "3lzmbrev00092") + require.NoError(t, world.dispatcher.HandleEvent(ctx, inB.create(t)), + "their comment in B must still federate: a ban is one community's ruling about its "+ + "own space, and a gate that reads 'is this author banned anywhere' silences them "+ + "everywhere on one moderator's decision") + assert.Equal(t, 1, outboundRowsFor(t, h.db, inB.atURI())) + assert.Greater(t, rowCount(t, h.db, "outbound_deliveries"), deliveriesBefore, + "and it really went out") +} + +// TestABannedAuthorsVoteIsRefusedInThatCommunityOnly is HIGH-1 for votes. +// +// A vote is the smallest thing an excluded author can still send, and the one +// they will send most: a banned user scrolling their own feed upvotes as they +// read. Every one of those is a delivery to an instance that refuses it. +func TestABannedAuthorsVoteIsRefusedInThatCommunityOnly(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + admitPost(t, world, mtAuthorDID, mbPostInBRKey, world.communityBDID, "3lzmbrev00100", 1_775_000_014_000_001) + postInB := mbPostATURI(mtAuthorDID, mbPostInBRKey) + + h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, nil) + _, banned := banFor(t, h.db, world.communityADID, mtAuthorDID) + require.True(t, banned, "precondition: banned in A, not in B") + + votesBefore := rowCount(t, h.db, "outbound_votes") + deliveriesBefore := rowCount(t, h.db, "outbound_deliveries") + + err := world.dispatcher.HandleEvent(ctx, + mbVoteEvent(t, mtAuthorDID, "3lzmbvote00001", "3lzmbrev00101", mtPostATURI, "up", 1_775_000_014_000_002)) + require.Error(t, err, + "a banned author's VOTE must be refused too: it is the highest-volume thing they can "+ + "still send into a community that refuses all of it") + assert.True(t, stderrors.Is(err, consume.ErrPermanentEvent), + "permanently, like every other refusal in this family (err=%v)", err) + assert.Contains(t, err.Error(), "author-banned", "with the cause named") + assert.Equal(t, votesBefore, rowCount(t, h.db, "outbound_votes"), + "and NO outbound vote state is written: that row is what a later Undo is rebuilt "+ + "from, so recording one for a vote we never sent leaves a retraction with nothing "+ + "behind it") + assert.Equal(t, deliveriesBefore, rowCount(t, h.db, "outbound_deliveries")) + + require.NoError(t, world.dispatcher.HandleEvent(ctx, + mbVoteEvent(t, mtAuthorDID, "3lzmbvote00002", "3lzmbrev00102", postInB, "up", 1_775_000_014_000_003)), + "their vote in B still counts: the ban is scoped to A's space") + assert.Greater(t, rowCount(t, h.db, "outbound_votes"), votesBefore, "and it was recorded") +} + +// TestABannedAuthorEditingAnAcceptedPostRemovesNothing is HIGH-2, and it is the +// destructive one. +// +// An edit re-runs admission. Once the ban gate rejects it, the post takes the +// "was accepted and now fails re-admission" path: the acceptance is DELETED, a +// removal is written, and a Delete{Page} is ENQUEUED to the community — which +// with removeData=false is content Lemmy explicitly KEPT. So a typo fix, minutes +// after a ban, deletes the author's post from the community view and pushes a +// Delete at the moderators who chose to leave it up. +// +// That is 17c-1's F2 in mirror image: there the engine reversed a moderator's +// removal on an edit; here it removes what a moderator preserved. And the +// removal it writes carries admission-revoked — the REVERSIBLE code, whose whole +// meaning is "our own decision, a corrective edit may undo it" — for a cause +// that no edit can ever satisfy. +func TestABannedAuthorEditingAnAcceptedPostRemovesNothing(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + _, _, err := h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + require.NoError(t, err, "precondition: the post is accepted in A") + + h.announceBlock(world.groupA, mbBlockActivity, mtAuthorDID, groupID, nil) + activitiesBefore := rowCount(t, h.db, "outbound_activities") + deliveriesBefore := rowCount(t, h.db, "outbound_deliveries") + + // The ordinary thing an author does minutes after being banned: fix a typo. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + mtPostEvent(t, "update", mtEditRev, mtEditCID, mtEditTime))) + + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + assert.NoError(t, err, + "the ACCEPTANCE MUST STAND: the moderators banned the author and left this post up — "+ + "removing it on their next edit destroys content the community chose to keep, and "+ + "the post disappearing from Coves is the visible half (err=%v)", err) + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionRemoval, world.digestRKey) + assert.True(t, errors.IsNotFound(err), + "and NO removal may be written — least of all under admission-revoked, whose meaning "+ + "is 'we withdrew this and a corrective edit may restore it', for a cause no edit "+ + "can ever satisfy (err=%v)", err) + + assert.Equal(t, activitiesBefore, rowCount(t, h.db, "outbound_activities"), + "NOTHING may be enqueued: a Delete{Page} here is the unrecoverable half — it is on "+ + "the wire, aimed at the moderators who deliberately kept this post, and no later "+ + "fix retracts it") + assert.Equal(t, deliveriesBefore, rowCount(t, h.db, "outbound_deliveries"), "...and no delivery") + + status, code := admissionFor(t, h.db, world.communityADID, mtPostATURI) + assert.Equal(t, "author-banned", code, + "and the ledger records WHY the edit did nothing: it is the surface an operator reads "+ + "when the author asks, and 'the edit silently vanished' is the only alternative") + assert.NotEqual(t, accept.StatusRemoved, status, + "while the post itself is not recorded as removed — it is still up, and the ledger "+ + "must not say otherwise") +} + +// TestAnAnnouncedBlockNamingAFediverseUserIsNotOurs pins the branch that decides +// a Block is somebody else's business. +// +// Lemmy announces every ban it issues, including bans of its OWN users, and we +// follow those communities — so this arrives constantly. It is enforced entirely +// on their instance; we hold no state that could apply it and no persona it +// could be about. The counter is what separates "we saw it and it was not ours" +// from "we never received it", which is the question asked the day the ban path +// looks broken. +func TestAnAnnouncedBlockNamingAFediverseUserIsNotOurs(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + before := metricValue(mbForeignSubjectMetric) + + const foreign = "https://lemmy.world/activities/announce/block/mb-foreign" + require.Equal(t, http.StatusAccepted, h.deliver(world.groupA, map[string]any{ + "id": foreign, + "type": "Announce", + "actor": groupID, + "audience": groupID, + "cc": []any{groupID + "/followers"}, + "object": map[string]any{ + "id": foreign + "/block", + "type": "Block", + "actor": modActorID, + "object": "https://lemmy.world/u/SomeLemmyUser", + "target": groupID, + "audience": groupID, + "cc": []any{groupID}, + }, + })) + h.drain() + + assert.Zero(t, bansIn(t, h.db, world.communityADID), + "a ban on a LEMMY user records nothing here: we hold no persona it could be about, "+ + "and inventing a row keyed on an id we cannot resolve to a DID would gate "+ + "admission on a subject that can never post") + assert.Equal(t, before+1, metricValue(mbForeignSubjectMetric), + "and it is counted as not-ours: every community we follow announces its own bans, so "+ + "this is the common case, and the day the ban path looks broken this counter is "+ + "what distinguishes 'none of them were ours' from 'we stopped receiving them'") + + event, err := h.events.GetEvent(ctx, foreign) + require.NoError(t, err) + assert.NotNil(t, event.ProcessedAt, "decided once, not retried: %s", event.Error) + assert.Nil(t, event.FailedAt, "nor poisoned") +} + +// mbCommentBy is one native comment by an author, replying to a post. +func mbCommentBy(did, rkey, parentURI, parentCID, rev string) nativeComment { + return nativeComment{ + did: did, rkey: rkey, + root: nativeRef{parentURI, parentCID}, + parent: nativeRef{parentURI, parentCID}, + createRev: rev, createCID: "bafyreih5xbmigkq5ikyhqiqhqzbwuqjxeitgtzwyxvjhfsfvswsxmnnf8a", + editRev: rev + "e", editCID: "bafyreih5xbmigkq5ikyhqiqhqzbwuqjxeitgtzwyxvjhfsfvswsxmnnf8b", + timeUS: 1_775_000_015_000_000, + } +} + +// mbVoteEvent builds a native vote commit. +func mbVoteEvent(t *testing.T, did, rkey, rev, subjectURI, direction string, timeUS int64) *consume.JetstreamEvent { + t.Helper() + frame := fmt.Sprintf(`{ + "did": %q, "time_us": %d, "kind": "commit", + "commit": { + "rev": %q, "operation": "create", + "collection": "social.coves.feed.vote", + "rkey": %q, "cid": %q, + "record": { + "$type": "social.coves.feed.vote", + "subject": {"uri": %q, "cid": %q}, + "direction": %q, + "createdAt": "2026-08-13T12:00:00.000Z" + } + } +}`, did, timeUS, rev, rkey, mtPostCID, subjectURI, mtPostCID, direction) + var event consume.JetstreamEvent + require.NoError(t, json.Unmarshal([]byte(frame), &event), "the frame must be valid wire JSON") + return &event +} + +// bansIn counts the bans standing in one community. +func bansIn(t *testing.T, db *sql.DB, communityDID string) int { + t.Helper() + var n int + require.NoError(t, db.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM community_bans WHERE community_did = $1`, communityDID).Scan(&n)) + return n +} diff --git a/internal/store/community_bans.go b/internal/store/community_bans.go index d3f63ce..c546175 100644 --- a/internal/store/community_bans.go +++ b/internal/store/community_bans.go @@ -4,6 +4,7 @@ import ( "context" "database/sql" "fmt" + "time" "tidepool/internal/errors" ) @@ -67,9 +68,21 @@ func (r *postgresCommunityBans) Ban(ctx context.Context, ban CommunityBan) (canc return 0, fmt.Errorf("ban %q in %q: %w", ban.SubjectDID, ban.CommunityDID, err) } - cancelled, err = cancelPendingForActorInCommunity(ctx, tx, ban.SubjectDID, ban.CommunityAPID) - if err != nil { - return 0, err + // THE ROW IS RECORDED EITHER WAY; THE CANCELLATION IS NOT. A Block whose + // expiry has already passed when it reaches us — delayed, redelivered after + // an outage, replayed from a backfill — is a faithful record of a ban that is + // over, and storing it keeps the audit trail honest (and idempotent, since a + // later redelivery finds the same row). But it is not in force, so it must + // not cancel work by an author nobody is currently excluding: a cancelled + // delivery is never re-queued. + // + // The condition is the SAME predicate Standing() reads, kept here rather than + // at the call site so no caller can cancel on a ban that does not apply. + if ban.ExpiresAt == nil || ban.ExpiresAt.After(time.Now()) { + cancelled, err = cancelPendingForActorInCommunity(ctx, tx, ban.SubjectDID, ban.CommunityAPID) + if err != nil { + return 0, err + } } if err := tx.Commit(); err != nil { return 0, fmt.Errorf("ban %q in %q: commit: %w", ban.SubjectDID, ban.CommunityDID, err) @@ -103,6 +116,17 @@ func (r *postgresCommunityBans) Lift(ctx context.Context, communityDID, subjectD } func (r *postgresCommunityBans) Standing(ctx context.Context, communityDID, subjectDID string) (bool, error) { + return standingBan(ctx, r.db, communityDID, subjectDID) +} + +func (r *postgresCommunityBans) StandingTx(ctx context.Context, tx *sql.Tx, communityDID, subjectDID string) (bool, error) { + if tx == nil { + return false, errors.NewValidationError("tx", "must not be nil") + } + return standingBan(ctx, tx, communityDID, subjectDID) +} + +func standingBan(ctx context.Context, q queryRower, communityDID, subjectDID string) (bool, error) { if communityDID == "" || subjectDID == "" { return false, errors.NewValidationError("ban", "community_did and subject_did must not be empty") } @@ -110,7 +134,7 @@ func (r *postgresCommunityBans) Standing(ctx context.Context, communityDID, subj // The expiry test is in the STATEMENT, not in Go, so no reader can forget // it: a lapsed ban is indistinguishable from no ban, and the only signal // that it lapsed is the clock — Lemmy sends nothing. - if err := r.db.QueryRowContext(ctx, ` + if err := q.QueryRowContext(ctx, ` SELECT EXISTS ( SELECT 1 FROM community_bans WHERE community_did = $1 AND subject_did = $2 diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index 0953ffa..ac5b9f6 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -247,6 +247,11 @@ type CommunityBans interface { // Re-banning preserves the original banned_at (a re-delivered Block is the // same ban twice) while taking the expiry, reason and removeData from the // new activity, which are the parts a moderator can genuinely re-issue. + // + // A ban whose expiry has ALREADY PASSED is still RECORDED — it is a faithful + // account of what the moderator sent, and a redelivery must find the same row + // — but it cancels nothing, because it is not in force and a cancelled + // delivery is never re-queued. Ban(ctx context.Context, ban CommunityBan) (cancelled int64, err error) // Lift removes the ban (Undo{Block}), reporting whether one was standing. @@ -259,6 +264,15 @@ type CommunityBans interface { // no ban: Lemmy sends no Undo when a timed ban runs out, so the clock is the // only thing that ever lifts it. Standing(ctx context.Context, communityDID, subjectDID string) (bool, error) + + // StandingTx is Standing on an existing transaction — the seam the + // acceptance commit uses to re-ask the question INSIDE the transaction that + // writes the acceptance and enqueues the delivery. A ban read before that + // transaction opens is a decision made about state that can change before it + // is acted on: the ban's own transaction cancels every pending delivery, so + // a post admitted after it lands is one the cancellation could never catch. + // A nil tx is an error satisfying errors.IsValidation. + StandingTx(ctx context.Context, tx *sql.Tx, communityDID, subjectDID string) (bool, error) } // Communities tracks the AP groups the bridge subscribes to and their