From 40ca3c574ab8db88354a6b6cf9c740a81d7549ee Mon Sep 17 00:00:00 2001 From: Bretton Date: Fri, 14 Aug 2026 14:30:24 -0700 Subject: [PATCH] fix(moderation): a lapsed ban acts on nothing, and the gate covers every write (17c-3 review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Two review streams. The ban's mechanism was sound; its edges were not. A LAPSED BAN WAS FULLY DESTRUCTIVE. applyBan never asked whether the ban it was applying had already expired, so a Block carrying a past `expires` was written, its pending deliveries cancelled, and with removeData its accepted posts purged under a terminal code — while Standing() correctly reported it inactive. The read path honoured the clock and the write path did not, and neither effect is recoverable. Expiry is now decided before any consequence. The ROW is still recorded — a faithful account of what the moderator sent, idempotent under redelivery — and the effects are gated: the cancellation inside Ban() beside the statement it guards, under the same predicate Standing() reads, so no caller can cancel on a ban that is not in force. AN UNREADABLE EXPIRY BECAME PERMANENT. Absent, present-but-unparseable and zero all collapsed to nil, and nil means forever — on the one field whose misreading cannot be corrected, since Lemmy sends no Undo when a temporary ban lapses. The three wire facts now fork: absent is permanent, parsed is weighed against the clock, unparseable is refused and persists nothing. Both spellings are read — `expires` (0.19) and AS2 `endTime` (newer Lemmy) — kept as they arrived. THE GATE COVERED POSTS ONLY. accept.decide never sees a comment or a vote, so a banned author's replies and votes kept enqueueing on the banning community's ordering key, to be refused by Lemmy and retried into poison with the cause three tables away. That the ban's own cancellation already swept that traffic — it joins actor and ordering key, not activity kind — is what showed this was a gap rather than a decision. Both paths now refuse with the same permanence discipline as the thread lock. A BANNED AUTHOR'S EDIT REMOVED CONTENT LEMMY KEPT. An edit of an already-accepted post took removeAccepted: acceptance deleted, removal written, Delete{Page} enqueued at the moderators who banned them — with removeData=false meaning Lemmy had explicitly kept that content — and the removal carried the REVERSIBLE admission-revoked code for a decision documented to survive an unban. It now records the ledger rejection and leaves the acceptance and the queue alone. The read side no longer fails open: NewEngine refuses to build without a ban store, matching the write side, because a gate that is absent does not fail — it admits. Also: a direct Undo{Block} is counted like a direct Block, so unban traffic is not invisible; the moderator's reason is persisted rather than leaving one audit surface quoting them and the other blank; and community_ap_id's comment now names its only real reader instead of one that does not exist. The in-transaction re-check narrows the admission race rather than closing it, and a delivery claimed just before the ban can still be POSTed — cancellation reaches the row, not the socket. Both are in FOLLOWUPS for 17e. Test-side: community_bans needed the same truncate line in a fourth and fifth package, and the failure never names it — a leftover row refuses the next test's author and reads as the engine being broken. Co-Authored-By: Claude Opus 5 (1M context) --- FOLLOWUPS.md | 40 ++ internal/accept/engine.go | 104 ++- internal/accept/outer_acceptance_test.go | 9 +- internal/accept/terminality_race_test.go | 128 ++++ internal/ap/vocab.go | 33 +- internal/consume/comments.go | 43 ++ internal/consume/dispatch.go | 7 + internal/consume/dispatch_test.go | 8 +- internal/consume/outer_acceptance_test.go | 8 +- internal/consume/votes.go | 4 + internal/db/migrations/027_community_bans.sql | 20 +- internal/ingest/consent.go | 9 + internal/ingest/handler.go | 4 +- internal/ingest/moderation.go | 76 +- internal/ingest/moderation_ban_test.go | 651 +++++++++++++++++- internal/store/community_bans.go | 32 +- internal/store/interfaces.go | 14 + 17 files changed, 1117 insertions(+), 73 deletions(-) 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 -- 2.51.2