diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index fd6fe1e..6ace7fd 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -306,13 +306,16 @@ func run(logger *slog.Logger) error { } materializer, err := materialize.New(materialize.Options{ - Fetcher: apClient, - Objects: objects, - Actors: actors, - Communities: communities, - Repos: repoManager, - Minter: mintGate, - Votes: voteAggregator, + Fetcher: apClient, + Objects: objects, + Actors: actors, + Communities: communities, + Repos: repoManager, + Minter: mintGate, + Votes: voteAggregator, + // Restoring a NATIVE post reads its pinned CID from here: the author's + // repo is not one this bridge hosts. + OutboundObjects: store.NewOutboundObjects(database), ServiceDID: serviceDID, ProfileRefreshTTL: cfg.ProfileRefreshTTL, MaxBlobBytes: cfg.MaxBlobBytes, diff --git a/internal/accept/engine.go b/internal/accept/engine.go index 6611589..7ef26ba 100644 --- a/internal/accept/engine.go +++ b/internal/accept/engine.go @@ -70,6 +70,17 @@ const ( // for federation purposes, so we must not sign an acceptance or deliver into // it. SECURITY: a communities row alone (existence) is not authority to bridge. DecisionCommunityNotFollowed = "community-not-followed" + // DecisionModeratorRemoved: the author edited a post the COMMUNITY has + // removed. The edit is not rejected — the author's record is theirs and it + // stands — but it does not re-enter a community that removed it, and no + // acceptance is written and nothing is enqueued. + // + // It is deliberately NOT RemovalCodeAdmissionRevoked: that code means "we + // withdrew this and a corrective edit may restore it", which is precisely + // the reasoning that must not reach a moderator's decision. The two must + // stay distinguishable in the ledger, because they are the two branches of + // what an edit against a standing removal is allowed to do. + DecisionModeratorRemoved = "moderator-removed" ) // RemovalCodeAdmissionRevoked is the removal `code` written when a post that WAS @@ -480,7 +491,10 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, ATURI: postURI, ID: consume.ActivityID(e.userOrigin, postURI, op, stored.LastActivitySeq), CommunityAPID: community.APGroupID, - Snapshot: snapshot, + // The mapping the enqueuer writes carries this binding; without it + // no announced moderation of this post can ever be authorized. + CommunityDID: communityDID, + Snapshot: snapshot, } // A post has no causal parent, so orderingKey is the author DID and there // is no parentATURI. @@ -502,17 +516,7 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, _, err = acceptrec.AcceptSubject(ctx, e.repos, communityDID, postURI, commit.CID, publishedAtOf(commit.Record), sideEffect) if stderrors.Is(err, acceptrec.ErrRemovalStands) { - // An admission-revoked removal stands (a prior failing edit withdrew the - // post), but admission passes NOW. The revocation was OUR decision, so a - // corrective edit AUTO-RESTORES rather than erroring and redriving forever - // against the terminal removal: delete the removal + write a fresh - // acceptance + enqueue (a Create, since Lemmy's live copy is gone). The - // same side effect rides the restore commit. - if _, rerr := acceptrec.Restore(ctx, e.repos, communityDID, postURI, commit.CID, - publishedAtOf(commit.Record), sideEffect); rerr != nil { - return fmt.Errorf("accept: restore %s into %s: %w", postURI, communityDID, rerr) - } - return nil + return e.editAgainstRemoval(ctx, did, communityDID, postURI, commit, sideEffect) } if err != nil { return fmt.Errorf("accept: admit %s into %s: %w", postURI, communityDID, err) @@ -520,6 +524,73 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, return nil } +// editAgainstRemoval decides what an edit may do when a removal already stands +// at the subject's rkey. WHOSE removal it is decides, and nothing else. +// +// - admission-revoked is OUR OWN decision: a prior edit failed admission and +// we withdrew the post. A corrective edit is exactly the event that should +// reverse it, so it auto-restores — delete the removal, write a fresh +// acceptance, enqueue (a Create, since Lemmy's live copy is gone), all on +// the one restore commit. +// - ANY OTHER CODE IS A MODERATOR'S DECISION AND IS TERMINAL. Reversing it +// would delete the moderators' removal record, write an acceptance over it, +// and — because the same side effect rides that commit — ENQUEUE the post +// back to the community that removed it. Every other consequence of getting +// this wrong is internal and correctable; that one is on the wire, at the +// people who made the decision. +// +// A terminal removal is recorded in the ledger and nothing else happens: no +// commit, no enqueue, no error. The author's edit is not a failure — their +// record is theirs and it stands — it simply does not re-enter a community that +// has removed it, and the ledger is where an operator reads why. +func (e *Engine) editAgainstRemoval(ctx context.Context, did, communityDID, postURI string, + commit *consume.CommitEvent, sideEffect repo.TxSideEffect) error { + + code, err := e.standingRemovalCode(ctx, communityDID, postURI) + if err != nil { + return err + } + if code != RemovalCodeAdmissionRevoked { + e.logger.Info("edit against a standing moderator removal: acceptance refused, nothing enqueued", + "community_did", communityDID, "post", postURI, "removal_code", code) + return e.admissions.Record(ctx, Admission{ + AuthorDID: did, + CommunityDID: communityDID, + PostURI: postURI, + Status: StatusRemoved, + DecisionCode: DecisionModeratorRemoved, + EvaluatedCID: commit.CID, + EvaluatedSnapshot: e.evaluatedSnapshot(commit), + }) + } + if _, rerr := acceptrec.Restore(ctx, e.repos, communityDID, postURI, commit.CID, + publishedAtOf(commit.Record), sideEffect); rerr != nil { + return fmt.Errorf("accept: restore %s into %s: %w", postURI, communityDID, rerr) + } + return nil +} + +// standingRemovalCode reads the `code` off the removal AcceptSubject refused +// against. A removal that has vanished between the refusal and this read is a +// genuine race — the moderators restored the post in the window — and returns +// an error so the event RETRIES: the retry's AcceptSubject finds no removal and +// admits the edit normally. Treating the miss as "not ours, terminal" would +// strand a post whose removal no longer exists. +// +// An unreadable code is treated as a moderator's, i.e. terminal. The direction +// is deliberate: the recoverable mistake is refusing an edit, and the +// unrecoverable one is pushing a removed post back at its community. +func (e *Engine) standingRemovalCode(ctx context.Context, communityDID, postURI string) (string, error) { + rkey := acceptrec.SubjectRKey(postURI) + record, _, err := e.repos.GetRecord(ctx, communityDID, acceptrec.CollectionRemoval, rkey) + if err != nil { + return "", fmt.Errorf("accept: read standing removal %s/%s/%s: %w", + communityDID, acceptrec.CollectionRemoval, rkey, err) + } + code, _ := record["code"].(string) + return code, nil +} + // removeAccepted withdraws a post that WAS accepted and now fails re-admission: // the acceptance is deleted and a removal (code admission-revoked) written in ONE // commit, carrying the Delete{Page} enqueue as the side effect. The admissions diff --git a/internal/consume/comments.go b/internal/consume/comments.go index 6f56823..a01f836 100644 --- a/internal/consume/comments.go +++ b/internal/consume/comments.go @@ -168,8 +168,12 @@ func (d *Dispatcher) enqueueComment(ctx context.Context, tx *sql.Tx, did, operat ATURI: stored.ATURI, ID: ActivityID(d.userOrigin, stored.ATURI, operation, stored.LastActivitySeq), CommunityAPID: stored.CommunityAPID, - ParentAPID: parentAPID, - Snapshot: stored.TranslatedSnapshot, + // Read off the stored outbound row rather than re-resolved: it is the + // same community this comment was admitted into, and the mapping the + // enqueuer writes needs it to be moderatable. + CommunityDID: stored.CommunityDID, + ParentAPID: parentAPID, + Snapshot: stored.TranslatedSnapshot, } // parentATURI carries the causal dependency (decision 15): delivery must // not present a reply to a peer before the thing it replies to. On a diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index 3f9b0ac..619710b 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -43,6 +43,12 @@ type CommentIntent struct { ID string // CommunityAPID is the target community's AP Group id. CommunityAPID string + // CommunityDID is that same community's bridged repo DID. It rides the + // intent so the enqueuer can BIND the mapping it writes to a community + // without growing a Communities dependency: every caller already holds the + // resolved community, and the binding is what later authorizes announced + // moderation of this object (materialize.CommunityDIDOf). + CommunityDID string // ParentAPID is the AP object id of the thing replied to, resolved through // ap_objects OR the parent's own outbound state (a native accepted post or // an earlier native comment, which have no ap_objects mapping). @@ -96,6 +102,9 @@ type PostIntent struct { // CommunityAPID is the target community's AP Group id — for a Page it goes // in `to` (the Page/Note addressing split), not `cc`. CommunityAPID string + // CommunityDID is that same community's bridged repo DID; see + // CommentIntent.CommunityDID for why it travels on the intent. + CommunityDID string // Snapshot is the translated state a Delete is rebuilt from and a // Create/Update{Page} is rendered from — the postv2 record plus resolved // context, same envelope shape as a comment's snapshot. diff --git a/internal/consume/subjects.go b/internal/consume/subjects.go index dad9d58..175948d 100644 --- a/internal/consume/subjects.go +++ b/internal/consume/subjects.go @@ -69,10 +69,20 @@ func (d *Dispatcher) resolveSubject(ctx context.Context, atURI string) (*resolve if err != nil { return nil, fmt.Errorf("resolve community %s: %w", communityDID, err) } - depth, err := d.recordedDepth(ctx, atURI) + depth, tombstoned, err := d.recordedState(ctx, atURI) if err != nil { return nil, err } + if tombstoned { + // The bridge already withdrew this object. This branch USED to + // be unreachable for native content — with no community_did on + // a bridge-origin mapping the lookup fell through to the + // outbound_objects branch below, whose IsTombstoned check + // caught it. Populating that column (17c) moves native parents + // up here, and without this check replies to and votes on an + // author-DELETED native post would start federating again. + return nil, nil + } return &resolvedSubject{ ATURI: atURI, APID: mapping.APID, @@ -129,22 +139,31 @@ func (d *Dispatcher) subjectCommunityDID(ctx context.Context, mapping *store.APO return communityDID, nil } -// recordedDepth reads a subject's own reply depth, which exists only if the -// bridge federated it OUTWARD too. A mapped subject with no outbound row is a -// post, or a Lemmy object whose depth this bridge does not track, so it counts -// as the top: replies to it are depth 1 (this returns 0). +// recordedState reads what the bridge's OWN outbound row says about a mapped +// subject: its reply depth, and whether it has been withdrawn. +// +// depth exists only if the bridge federated the subject outward too. A mapped +// subject with no outbound row is a post, or a Lemmy object whose depth this +// bridge does not track, so it counts as the top: replies to it are depth 1 +// (this returns 0). +// +// tombstoned is the same fact the outbound_objects branch of resolveSubject +// checks, read HERE because both facts come off one row and the caller needs +// them together — a second read could see a delete land between them and +// federate against a subject the first read had already shown as live. // -// Only a genuine MISS defaults to 0. A real error — postgres down, a timeout — -// PROPAGATES: recording depth 0 off a dropped connection would federate a -// deeply nested comment at the wrong nesting and, past Lemmy's cap, keep -// federating ones it will reject, silently and with no retry. -func (d *Dispatcher) recordedDepth(ctx context.Context, atURI string) (int, error) { +// Only a genuine MISS defaults to zero values. A real error — postgres down, a +// timeout — PROPAGATES: reading depth 0 off a dropped connection would federate +// a deeply nested comment at the wrong nesting and, past Lemmy's cap, keep +// federating ones it will reject, silently and with no retry; and reading +// "not tombstoned" off one would resurrect deleted content. +func (d *Dispatcher) recordedState(ctx context.Context, atURI string) (depth int, tombstoned bool, err error) { state, err := d.objects.GetByATURI(ctx, atURI) if errors.IsNotFound(err) { - return 0, nil + return 0, false, nil } if err != nil { - return 0, fmt.Errorf("read recorded depth for %s: %w", atURI, err) + return 0, false, fmt.Errorf("read recorded state for %s: %w", atURI, err) } - return state.Depth, nil + return state.Depth, state.IsTombstoned(), nil } diff --git a/internal/db/migrations/024_ap_object_community_backfill.sql b/internal/db/migrations/024_ap_object_community_backfill.sql new file mode 100644 index 0000000..7d379b0 --- /dev/null +++ b/internal/db/migrations/024_ap_object_community_backfill.sql @@ -0,0 +1,33 @@ +-- +goose Up +-- Task 17c-1: bind bridge-origin objects to the community they were federated +-- into, so announced moderation of NATIVE content can be authorized at all. +-- +-- materialize.CommunityDIDOf is the bridge-wide answer to "whose content is +-- this?", and for an author-owned postv2 it cannot be derived: the record lives +-- in the AUTHOR's repo, so `did` names the author, not the community. With the +-- column unset it answers "" and authorizeAnnouncedContentDelete refuses every +-- announced removal — the accident that has been standing in for authorization. +-- The enqueuer now carries community_did on the mapping it writes; this +-- backfills the rows written before it did. +-- +-- The backfill is an EXACT JOIN, not a derivation. outbound_objects is the +-- bridge's own record of what it federated where, keyed by the same at-uri, and +-- it already carries community_did. Guessing (say, from the acceptance records +-- in each followed community) could bind a post to the wrong community, and a +-- wrong binding is worse than none: it would let one community's moderators +-- remove another community's content. +UPDATE ap_objects a + SET community_did = o.community_did + FROM outbound_objects o + WHERE a.at_uri = o.at_uri + AND a.origin = 'bridge' + AND a.community_did IS NULL + AND o.community_did <> ''; + +-- +goose Down +-- Deliberately NOT reversed. Clearing community_did would silently re-disable +-- moderation of native content — the refusal looks identical to "no moderator +-- acted" — and the values are recoverable from outbound_objects at any time by +-- re-running the statement above. A Down that re-creates a security-relevant +-- blind spot is worse than one that leaves a correct binding in place. +SELECT 1; diff --git a/internal/ingest/echo_delete_test.go b/internal/ingest/echo_delete_test.go index e0a5584..db9443c 100644 --- a/internal/ingest/echo_delete_test.go +++ b/internal/ingest/echo_delete_test.go @@ -137,8 +137,10 @@ func setupNativeDeleteEcho(t *testing.T, h *harness, communityDID, authorDID str require.NoError(t, err, "federating the post must have written its bridge-origin mapping") require.Equal(t, store.OriginBridge, mapping.Origin) require.Equal(t, materialize.CollectionPostV2, mapping.Collection) - require.Empty(t, mapping.CommunityDID, - "precondition: today's enqueuer records no community on the mapping — the accident this test exists to outlive") + // (The tripwire that stood here — "today's enqueuer records no community" — + // fired in 17c-1 and has been removed. Its whole purpose was to fail on the + // day M1b stopped simulating: community_did is populated now, so this + // scenario is the real world rather than a construction of it.) mapping.CommunityDID = communityDID mapping.AuthorDID = authorDID _, err = h.objects.PutMapping(ctx, *mapping) diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index e4e4628..4101f31 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -249,7 +249,15 @@ func newHarness(t *testing.T) *harness { // persona from another package's run would make an echo look like // someone else's. "outbound_deliveries", "outbound_activities", "outbound_objects", - "outbound_votes", "ap_actors", "vote_events", "vote_aggregates") + "outbound_votes", "ap_actors", "vote_events", "vote_aggregates", + // Consumer + engine state, for the tests that drive a real Jetstream + // commit through the acceptance engine. The rev gate is the one that + // bites: a leftover jetstream_record_revs row makes the SECOND test to + // use a given rev skip its own fixture as already-applied, and the + // failure surfaces as a missing acceptance record rather than as + // anything about revs. + "jetstream_record_revs", "jetstream_dead_letters", "consumer_cursors", + "admissions", "federation_prefs") custodian, err := identity.NewCustodian(testKEK) require.NoError(t, err) diff --git a/internal/ingest/moderation_terminal_test.go b/internal/ingest/moderation_terminal_test.go new file mode 100644 index 0000000..db77741 --- /dev/null +++ b/internal/ingest/moderation_terminal_test.go @@ -0,0 +1,342 @@ +package ingest + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/accept" + "tidepool/internal/consume" + "tidepool/internal/errors" + "tidepool/internal/materialize" + "tidepool/internal/outbound" + "tidepool/internal/personas" + "tidepool/internal/store" +) + +// TASK 17c-1 — A MODERATOR'S REMOVAL IS REAL, AND IT IS TERMINAL. +// +// The removal path itself is already built: materialize.RemovePost does the one +// multi-op commit and moderateAnnouncedDelete already calls it. What has kept it +// unreachable for NATIVE content is authorization — CommunityDIDOf returns "" +// for a bridge-origin mapping with no community_did, so an announced delete of a +// native post is refused before it can act. +// +// Populating that column is this sub-run's prerequisite, and it arms a loaded +// gun (F2): the engine's auto-restore does not read the standing removal's code. +// It was written for OUR decisions — an admission-revoked removal that a +// corrective edit SHOULD reverse. The moment a moderator-discretion removal can +// exist for a native post, the next author edit deletes the moderator's removal, +// writes a fresh acceptance, and enqueues Update{Page} BACK TO THE COMMUNITY +// THAT REMOVED IT. +// +// That is the failure this file exists to prevent, and it is unrecoverable in +// the way that matters: it is not a number that reads wrong internally, it is a +// moderation decision reversed and pushed outward, over the wire, at the +// moderators who made it. +// +// THE FIXTURE HAS TWO COMMUNITIES, CO-HOSTED. Decision 18's rule is a +// CONJUNCTION — the signer must BE the community AND the target must belong to +// it — and with one community those are the same fact, so an implementation +// checking either one passes. Lemmy hosts many communities per instance and +// SameAuthority is true across all of them, so B-moderates-A's-content is +// ordinary traffic, not a thought experiment. +const ( + mtUserOrigin = "https://coves.social" + + // Community A owns the post. + mtCommunityAName = "technology" + // Community B is a DIFFERENT community on the SAME instance. Same authority, + // same signer host, no claim whatsoever on A's content. + mtCommunityBName = "science" + mtCommunityBAPID = "https://lemmy.world/c/" + mtCommunityBName + + mtAuthorDID = "did:plc:mtnativeauthor0001" + mtAuthorHandle = "mtauthor.coves.social" + mtPostRKey = "3lzmtpost00001" + mtPostATURI = "at://" + mtAuthorDID + "/social.coves.community.postv2/" + mtPostRKey + mtPostAPID = mtUserOrigin + "/ap/object/" + mtAuthorDID + + "/social.coves.community.postv2/" + mtPostRKey + mtPostCID = "bafyreievgu2ty7qbiaaom5zhmkznsnajuzideek3lo7e65dwqlrvrxnmo4" + mtEditCID = "bafyreib2rxk3rybk3aobmv5cjuql3bm2twh4jo5uxgf5kpqrsqxi3jgxte" + mtPostRev = "3lzmtrev000001" + mtEditRev = "3lzmtrev000002" + mtTimeUS = int64(1_775_000_000_000_000) + mtEditTime = int64(1_775_000_000_000_001) +) + +// moderationWorld is a native post accepted into community A, with community B +// standing beside it on the same instance. +type moderationWorld struct { + groupA *remoteActor + groupB *remoteActor + communityADID string + communityBDID string + dispatcher *consume.Dispatcher + digestRKey string +} + +func newModerationWorld(t *testing.T, h *harness) moderationWorld { + t.Helper() + ctx := context.Background() + groupA := h.subscribeTechnology() + h.serveLemmyWorldContent() + communityADID := testDIDFor(mtCommunityAName, "lemmy.world") + + // --- Community B: followed, co-hosted, and able to sign its own deliveries. + communityBDID := testDIDFor(mtCommunityBName, "lemmy.world") + _, err := h.communities.UpsertCommunity(ctx, store.Community{ + APGroupID: mtCommunityBAPID, + DID: communityBDID, + PreferredUsername: mtCommunityBName, + Instance: "lemmy.world", + }) + require.NoError(t, err) + require.NoError(t, h.communities.SetFollowState(ctx, mtCommunityBAPID, store.FollowStateAccepted)) + groupB := h.newRemoteActor(mtCommunityBAPID, map[string]any{ + "type": "Group", + "id": mtCommunityBAPID, + "preferredUsername": mtCommunityBName, + "inbox": mtCommunityBAPID + "/inbox", + "published": "2024-01-01T00:00:00.000000Z", + }) + + // --- The native side: a persona service, the real enqueuer, the real + // acceptance engine, and the real consumer in front of them. + userOrigin, err := personas.New(personas.Options{ + DB: h.db, + Custodian: h.custodian, + UserOrigin: mtUserOrigin, + }) + require.NoError(t, err) + enqueuer, err := outbound.NewEnqueuer(outbound.EnqueuerOptions{ + DB: h.db, + Translator: outbound.NewTranslator(mtUserOrigin), + Inboxes: outbound.NewInboxResolver(h.client, time.Minute), + Actors: store.NewAPActors(h.db), + UserOrigin: mtUserOrigin, + }) + require.NoError(t, err) + engine, err := accept.NewEngine(accept.Options{ + Repos: h.manager, + Enqueuer: enqueuer, + Actors: userOrigin, + Resolver: mtResolver{}, + Communities: h.communities, + Objects: store.NewOutboundObjects(h.db), + Prefs: store.NewFederationPrefs(h.db), + Admissions: accept.NewAdmissions(h.db), + UserOrigin: mtUserOrigin, + }) + require.NoError(t, err) + dispatcher, err := consume.NewDispatcher(consume.Options{ + DB: h.db, + Actors: userOrigin, + Enqueuer: enqueuer, + Resolver: mtResolver{}, + Engine: engine, + UserOrigin: mtUserOrigin, + }) + require.NoError(t, err) + + // --- The post is admitted through the real path: acceptance record, + // outbound state, ap_objects mapping, one Create{Page} enqueued. + require.NoError(t, dispatcher.HandleEvent(ctx, mtPostEvent(t, "create", mtPostRev, mtPostCID, mtTimeUS))) + + digest := testDigestRKey(mtPostATURI) + _, _, err = h.manager.GetRecord(ctx, communityADID, materialize.CollectionAcceptance, digest) + require.NoError(t, err, "precondition: the post is accepted into community A") + + // --- The 17c PREREQUISITE, constructed the way 17c will leave the world. + // Without community_did on the mapping, CommunityDIDOf returns "" and + // every announced moderation action is refused before it is evaluated — + // which is the accident that has been standing in for authorization. + // GREEN carries this column through the intent; the fixture states the + // end state so these behaviours can be pinned against it. + mapping, err := h.objects.GetByAPID(ctx, mtPostAPID) + require.NoError(t, err, "the enqueuer maps the federated post") + require.Equal(t, store.OriginBridge, mapping.Origin) + mapping.CommunityDID = communityADID + _, err = h.objects.PutMapping(ctx, *mapping) + require.NoError(t, err) + + return moderationWorld{ + groupA: groupA, groupB: groupB, + communityADID: communityADID, communityBDID: communityBDID, + dispatcher: dispatcher, digestRKey: digest, + } +} + +// mtPostEvent builds a postv2 commit frame the consumer accepts. +func mtPostEvent(t *testing.T, operation, rev, cid string, timeUS int64) *consume.JetstreamEvent { + t.Helper() + frame := fmt.Sprintf(`{ + "did": %q, "time_us": %d, "kind": "commit", + "commit": { + "rev": %q, "operation": %q, + "collection": "social.coves.community.postv2", + "rkey": %q, "cid": %q, + "record": { + "$type": "social.coves.community.postv2", + "community": %q, + "title": "a native post the community moderates", + "content": "the body, later edited", + "createdAt": "2026-08-13T10:00:00.000Z" + } + } +}`, mtAuthorDID, timeUS, rev, operation, mtPostRKey, cid, testDIDFor(mtCommunityAName, "lemmy.world")) + var event consume.JetstreamEvent + require.NoError(t, json.Unmarshal([]byte(frame), &event), "the frame must be valid wire JSON") + return &event +} + +type mtResolver struct{} + +func (mtResolver) ResolveDIDHandle(context.Context, string) (string, error) { + return mtAuthorHandle, nil +} + +// admissionFor reads the ledger row for one (community, post). +func admissionFor(t *testing.T, db *sql.DB, communityDID, postURI string) (status, code string) { + t.Helper() + err := db.QueryRow(` + SELECT status, decision_code FROM admissions + WHERE community_did = $1 AND post_uri = $2`, communityDID, postURI).Scan(&status, &code) + if err == sql.ErrNoRows { + return "", "" + } + require.NoError(t, err) + return status, code +} + +// TestModeratorRemovalSurvivesAnAuthorEdit is the OUTER CONTRACT for 17c-1. +// +// GIVEN a native post federated into community A and a moderator removal +// standing against it, WHEN the author edits the post, THEN the removal +// survives, no acceptance is written, NOTHING goes outbound, and the ledger +// records the terminal decision. +// +// The outbound assertion is not a nicety. The harm is not that the removal +// disappeared from our repo — it is that we take the moderators' decision, +// reverse it, and PUSH THE POST BACK AT THEM. Every other consequence is +// internal and correctable; that one is on the wire. +func TestModeratorRemovalSurvivesAnAuthorEdit(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + // --- A moderator of community A removes the post: Delete WITH summary, + // announced by the community that owns it. + reason := "off topic for this community" + h.announceDeleteWithSummary(world.groupA, + "https://lemmy.world/activities/announce/delete/mt-removal", mtPostAPID, &reason) + + removal, _, err := h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionRemoval, world.digestRKey) + require.NoError(t, err, + "precondition: with community_did populated the removal path is REACHABLE — this is "+ + "the whole point of the column, and the rest of this test is its consequence") + require.Equal(t, "moderator-discretion", removal["code"]) + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + require.True(t, errors.IsNotFound(err), "precondition: the acceptance was withdrawn (err=%v)", err) + + activitiesBefore := rowCount(t, h.db, "outbound_activities") + deliveriesBefore := rowCount(t, h.db, "outbound_deliveries") + + // --- WHEN: the author edits their post. An ordinary commit, the kind that + // happens minutes later when someone fixes a typo. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + mtPostEvent(t, "update", mtEditRev, mtEditCID, mtEditTime))) + + // --- THEN: the moderators' decision stands. + // assert, not require: every consequence below is a separate harm, and the + // outbound one is the harm that leaves the building. + stillRemoved, _, err := h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionRemoval, world.digestRKey) + if assert.NoError(t, err, + "the removal MUST survive the edit: auto-restore exists to reverse OUR OWN "+ + "admission-revoked decisions, and applying it to a moderator's removal reverses "+ + "a decision we had no part in") { + assert.Equal(t, "moderator-discretion", stillRemoved["code"], + "and it must still be the moderator's removal, not one we rewrote") + assert.Equal(t, reason, stillRemoved["reason"], "with their reason intact") + } + + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + assert.True(t, errors.IsNotFound(err), + "and NO acceptance may be written: an acceptance beside a standing removal is a post "+ + "that reads as both admitted and removed, and Coves resolves that by showing it "+ + "(err=%v)", err) + + assert.Equal(t, activitiesBefore, rowCount(t, h.db, "outbound_activities"), + "NOTHING may be enqueued: this is the unrecoverable half — an Update{Page} here "+ + "pushes the removed post back at the very moderators who removed it, over the "+ + "wire, where no later fix can retract it") + assert.Equal(t, deliveriesBefore, rowCount(t, h.db, "outbound_deliveries"), + "...and no delivery either") + + // --- And the decision is legible afterwards: an operator asking why the + // edit did not publish must find the answer in the ledger. + status, code := admissionFor(t, h.db, world.communityADID, mtPostATURI) + assert.Equal(t, accept.StatusRemoved, status, + "the ledger records the post as REMOVED, not accepted: it is the surface an operator "+ + "reads when the author asks why their edit did nothing") + assert.NotEmpty(t, code, "with a machine-readable why") + assert.NotEqual(t, accept.RemovalCodeAdmissionRevoked, code, + "and it must NOT read as our own admission revocation — that code means 'we withdrew "+ + "this and a corrective edit may restore it', which is exactly the reasoning that "+ + "must not apply here. The specific string is GREEN's to choose") +} + +// TestCrossCommunityRemovalIsRefused is the RELATIONAL case: community B, on the +// same instance as A, announces a removal of A's post. +// +// Decision 18's conjunction collapses in a one-community fixture, so this is the +// only shape that can tell "the signer IS the community" from "the target is IN +// the community". Lemmy co-hosts communities by design and SameAuthority is true +// across all of them, so nothing about this delivery is malformed — it is simply +// not B's post to moderate. +func TestCrossCommunityRemovalIsRefused(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + reason := "not your community's post" + h.announceDeleteWithSummary(world.groupB, + "https://lemmy.world/activities/announce/delete/mt-cross", mtPostAPID, &reason) + + _, _, err := h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionRemoval, world.digestRKey) + assert.True(t, errors.IsNotFound(err), + "community B may not remove community A's post: one moderator team would otherwise "+ + "be able to withdraw content from every community co-hosted with it (err=%v)", err) + + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, world.digestRKey) + assert.NoError(t, err, + "and A's acceptance is untouched: the refusal must leave the post exactly as A "+ + "admitted it") + + _, _, err = h.manager.GetRecord(ctx, + world.communityBDID, materialize.CollectionRemoval, world.digestRKey) + assert.True(t, errors.IsNotFound(err), + "nor may B record the removal in its OWN repo: a removal for a post that was never "+ + "accepted there is a moderation record about somebody else's content (err=%v)", err) + + // The post is still live for its own community, so an edit still publishes. + 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, + "a refused cross-community removal must not leave the post in a state where its own "+ + "community's edits stop working") +} diff --git a/internal/materialize/acceptance.go b/internal/materialize/acceptance.go index ccd4785..b3972f8 100644 --- a/internal/materialize/acceptance.go +++ b/internal/materialize/acceptance.go @@ -191,16 +191,15 @@ func (m *Materializer) RestorePost(ctx context.Context, mapping *store.APObjectM return nil } - _, currentCID, err := m.repos.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + currentCID, ok, err := m.restorePin(ctx, mapping) if err != nil { - // No post to re-accept. The removal stays: it is terminal on the URI, - // and an acceptance pinning nothing would be worse than none. - if errors.IsNotFound(err) { - m.logger.Warn("restore target has no record to re-accept; leaving the removal in place", - "ap_id", mapping.APID, "at_uri", mapping.ATURI) - return nil - } - return fmt.Errorf("materialize: read %s for restore: %w", mapping.ATURI, err) + return err + } + if !ok { + // No version to pin. The removal stays: it is terminal on the URI, and + // an acceptance pinning nothing — or pinning a record that is gone — + // would be worse than none. + return nil } acceptance := map[string]any{ @@ -223,6 +222,67 @@ func (m *Materializer) RestorePost(ctx context.Context, mapping *store.APObjectM return nil } +// restorePin is the CID a fresh acceptance pins, dispatched on ORIGIN because +// the two origins keep the post's current version in different places. +// +// A fediverse-origin post was materialized into a repo this bridge hosts, so +// Repos can read it back. A BRIDGE-ORIGIN (native) post cannot: mapping.DID is +// the AUTHOR's DID and their PDS is not ours, so that read returns NotFound for +// every native post — which is why restoring one used to log "no record to +// re-accept" and leave every moderator restore silently unapplied. +// +// For those, outbound_objects.LastCID is the bridge's own record of the version +// it last federated. Migration 018 marks that column PROVENANCE ONLY, and this +// respects the reason: the caveat is about using it as an ORDERING GATE (a rev +// read from it is a check-then-write race). As a PIN it is exactly right — it +// names the version the community last saw — and it self-heals, because the +// author's next edit re-accepts against the fresh CID. +// +// ok=false means "no version to pin, leave the removal standing". A TOMBSTONED +// outbound row is the sharpest case: the author deleted the post while it was +// removed, so restoring would publish an acceptance for a record that no longer +// exists — the community asserting it admitted something deleted. +func (m *Materializer) restorePin(ctx context.Context, mapping *store.APObjectMapping) (string, bool, error) { + if mapping.Origin != store.OriginBridge { + _, cid, err := m.repos.GetRecord(ctx, mapping.DID, mapping.Collection, mapping.RKey) + if err != nil { + if errors.IsNotFound(err) { + m.logger.Warn("restore target has no record to re-accept; leaving the removal in place", + "ap_id", mapping.APID, "at_uri", mapping.ATURI) + return "", false, nil + } + return "", false, fmt.Errorf("materialize: read %s for restore: %w", mapping.ATURI, err) + } + return cid, true, nil + } + + if m.outbound == nil { + m.logger.Warn("restore of a native post needs outbound state, which is not wired; leaving the removal in place", + "ap_id", mapping.APID, "at_uri", mapping.ATURI) + return "", false, nil + } + state, err := m.outbound.GetByATURI(ctx, mapping.ATURI) + if err != nil { + if errors.IsNotFound(err) { + m.logger.Warn("restore target has no outbound state to re-accept; leaving the removal in place", + "ap_id", mapping.APID, "at_uri", mapping.ATURI) + return "", false, nil + } + return "", false, fmt.Errorf("materialize: read outbound state for %s: %w", mapping.ATURI, err) + } + if state.IsTombstoned() { + m.logger.Warn("restore refused: the author deleted this post while it was removed", + "ap_id", mapping.APID, "at_uri", mapping.ATURI) + return "", false, nil + } + if state.LastCID == "" { + m.logger.Warn("restore target has no recorded CID to pin; leaving the removal in place", + "ap_id", mapping.APID, "at_uri", mapping.ATURI) + return "", false, nil + } + return state.LastCID, true, nil +} + // moderationTarget resolves the community, subject uri and digest rkey a // moderation transition acts on, refusing anything that is not a postv2. // Moderation records are postv2-only: a pre-flip post has no acceptance to diff --git a/internal/materialize/materializer.go b/internal/materialize/materializer.go index d8095ee..7ce55eb 100644 --- a/internal/materialize/materializer.go +++ b/internal/materialize/materializer.go @@ -150,6 +150,12 @@ type Options struct { // Votes scrubs a deleted actor's vote_events rows alongside the record // scrub (optional; nil skips it). Votes VoteScrubber + // OutboundObjects is the bridge's own outbound state for NATIVE records. + // RestorePost needs it: a native post lives in the AUTHOR's repo, which + // this bridge does not host, so its CID cannot be read back through Repos. + // OPTIONAL: nil leaves a bridge-origin restore refusing (it logs and leaves + // the removal standing) exactly as it did before this seam existed. + OutboundObjects store.OutboundObjects // ServiceDID is the bridge's own DID: community.profile createdBy and // hostedBy (PLAN.md locked decision 6). ServiceDID string @@ -173,6 +179,7 @@ type Options struct { type Materializer struct { fetcher Fetcher objects store.APObjects + outbound store.OutboundObjects actors store.BridgedActors communities store.Communities repos *repo.Manager @@ -232,6 +239,7 @@ func New(opts Options) (*Materializer, error) { m := &Materializer{ fetcher: opts.Fetcher, objects: opts.Objects, + outbound: opts.OutboundObjects, actors: opts.Actors, communities: opts.Communities, repos: opts.Repos, diff --git a/internal/outbound/enqueuer.go b/internal/outbound/enqueuer.go index d63a878..8356451 100644 --- a/internal/outbound/enqueuer.go +++ b/internal/outbound/enqueuer.go @@ -165,7 +165,7 @@ func (e *Enqueuer) EnqueueActivity(ctx context.Context, tx *sql.Tx, actorDID, or // object; a self-delete's object was mapped on its create. ok=false means no // mapping is written. func (e *Enqueuer) objectMapping(intent consume.Intent) (store.APObjectMapping, bool, error) { - var atURI, apType string + var atURI, apType, communityDID string var snapshot []byte switch typed := intent.(type) { case consume.CommentIntent: @@ -173,11 +173,13 @@ func (e *Enqueuer) objectMapping(intent consume.Intent) (store.APObjectMapping, return store.APObjectMapping{}, false, nil } atURI, apType, snapshot = typed.ATURI, "Note", typed.Snapshot + communityDID = typed.CommunityDID case consume.PostIntent: if typed.Op == "delete" { return store.APObjectMapping{}, false, nil } atURI, apType, snapshot = typed.ATURI, "Page", typed.Snapshot + communityDID = typed.CommunityDID default: return store.APObjectMapping{}, false, nil } @@ -199,9 +201,15 @@ func (e *Enqueuer) objectMapping(intent consume.Intent) (store.APObjectMapping, OriginInstance: e.originHost, Origin: store.OriginBridge, DID: parts[0], - Collection: parts[1], - RKey: parts[2], - CID: cid, + // The COMMUNITY this object was federated into. Without it + // CommunityDIDOf answers "" for a bridge-origin mapping — an + // author-owned postv2 lives in the AUTHOR's repo, so the community + // cannot be read off DID — and every announced moderation action + // against native content is refused before it is even evaluated. + CommunityDID: communityDID, + Collection: parts[1], + RKey: parts[2], + CID: cid, }, true, nil } diff --git a/internal/store/ap_objects.go b/internal/store/ap_objects.go index 2bc0840..533ce41 100644 --- a/internal/store/ap_objects.go +++ b/internal/store/ap_objects.go @@ -56,7 +56,12 @@ func (r *postgresAPObjects) putMapping(ctx context.Context, q queryRower, mappin origin = EXCLUDED.origin, did = EXCLUDED.did, author_did = EXCLUDED.author_did, - community_did = EXCLUDED.community_did, + -- COALESCE, never a bare overwrite: community_did is the binding + -- that authorizes announced moderation of this object, and a + -- re-put that simply omits it (a re-materialization, a legacy + -- write path) would NULL a good binding and make moderation refuse + -- forever, silently. A write that HAS the value still wins. + community_did = COALESCE(EXCLUDED.community_did, ap_objects.community_did), collection = EXCLUDED.collection, rkey = EXCLUDED.rkey, at_uri = EXCLUDED.at_uri, diff --git a/internal/votes/reseed_faults_test.go b/internal/votes/reseed_faults_test.go index d9b7ab8..4bc4b76 100644 --- a/internal/votes/reseed_faults_test.go +++ b/internal/votes/reseed_faults_test.go @@ -95,6 +95,20 @@ func deliverTolerating(t *testing.T, l *lifecycle) { } } +// releaseHeldDelivery lets the worker pick the delivery up again. A settlement +// that fails after a confirmed POST HOLDS its delivery — still pending, still +// claimed for the rest of the lease — so a redelivery is a minute away in +// production and unreachable inside a test that refuses to sleep. Clearing the +// claim is how the lease lapsing is spelled here; what is under test is what +// happens when the worker gets its next look, not the clock that gives it one. +func releaseHeldDelivery(t *testing.T, l *lifecycle) { + t.Helper() + _, err := l.db.ExecContext(context.Background(), + `UPDATE outbound_deliveries SET claimed_until = NULL, next_attempt_at = now() + WHERE state = 'pending'`) + require.NoError(t, err) +} + // TestDeliveredLikeNeverStrandsThePendingLedgerRow is (a). // // The POST succeeded — Lemmy holds the vote — and the delivery is recorded as @@ -129,6 +143,7 @@ func TestDeliveredLikeNeverStrandsThePendingLedgerRow(t *testing.T) { // Whatever the mechanism, once the fault clears the worker must be able to // finish the job — a delivery it can no longer claim cannot be finished. faulty.failSet = false + releaseHeldDelivery(t, l) deliverTolerating(t, l) assert.Equal(t, string(store.DeliveredStateDelivered), l.state(t), "after the blip passes the ledger must catch up: the vote IS at Lemmy") @@ -179,6 +194,7 @@ func TestDeliveredUndoNeverStrandsTheDeliveredLedgerRow(t *testing.T) { } faulty.failDelete = false + releaseHeldDelivery(t, l) deliverTolerating(t, l) assert.Equal(t, "", l.state(t), "once the blip passes the row must go: the peer is not holding this vote")