diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index 675b626..7357381 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -70,8 +70,9 @@ task documents and git history rather than this list. index, but the outer `c.seq = ANY(ARRAY(...)) FOR UPDATE` re-check has no index on `seq` alone (`seq` is BIGSERIAL, not the PK `(activity_id, target_inbox)`). Unlike inbox_events (where `id` IS the PK), this is not a - point-fetch. A dedicated `UNIQUE INDEX (seq)` in its OWN migration (021 — - amending the already-applied 020 is a silent no-op under goose) would + point-fetch. A dedicated `UNIQUE INDEX (seq)` in its OWN migration (the next + free number — 023 as of this writing; amending an already-applied migration + is a silent no-op under goose) would restore O(keys × log N); add it if the claim path profiles hot. - **OutboundDeliveries is a 12-method interface** (go-proverbs SHOULD). It is one cohesive repository seam but the worker uses only the diff --git a/internal/accept/admin.go b/internal/accept/admin.go index 29860bf..12e7d83 100644 --- a/internal/accept/admin.go +++ b/internal/accept/admin.go @@ -171,7 +171,7 @@ func adminBearer(token string, logger *slog.Logger) func(http.Handler) http.Hand } } -// writeJSON is the shared JSON responder GREEN's handlers use. +// writeJSON is the shared JSON responder the admin handlers use. func writeJSON(w http.ResponseWriter, status int, body any) { w.Header().Set("Content-Type", "application/json; charset=utf-8") w.WriteHeader(status) diff --git a/internal/accept/engine.go b/internal/accept/engine.go index 3d3577d..6611589 100644 --- a/internal/accept/engine.go +++ b/internal/accept/engine.go @@ -22,6 +22,7 @@ import ( "log/slog" "strings" "time" + "unicode/utf8" "github.com/bluesky-social/indigo/atproto/atdata" "github.com/bluesky-social/indigo/atproto/lexicon" @@ -29,6 +30,7 @@ import ( "tidepool/internal/acceptrec" "tidepool/internal/consume" "tidepool/internal/errors" + "tidepool/internal/materialize" "tidepool/internal/repo" "tidepool/internal/store" "tidepool/lexicons" @@ -63,6 +65,11 @@ const ( // community's admissions row as a no-op annotation, if at all — the engine // writes nothing to the target community.) DecisionCommunityImmutable = "community-immutable" + // DecisionCommunityNotFollowed: the target community is not one Tidepool has + // an ACCEPTED Follow to (follow_state none/pending). We are not its key holder + // 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" ) // RemovalCodeAdmissionRevoked is the removal `code` written when a post that WAS @@ -116,6 +123,10 @@ type Options struct { // ADMISSION_MAX_PER_AUTHOR_PER_COMMUNITY). 0 means UNLIMITED (the generous // default); a positive value is the cap the rate check enforces. MaxPerAuthorPerCommunity int + // Now is the clock a moderation/removal record's createdAt is stamped from — + // the DECISION time, not the post's publication time (a removal on an old post + // is dated ~now). Injectable for tests. Nil uses time.Now. + Now func() time.Time // UserOrigin is AP_USER_ORIGIN: the origin every deterministic activity id // is minted under. UserOrigin string @@ -136,6 +147,7 @@ type Engine struct { apActors store.APActors maxPerCommunity int catalog *lexicon.BaseCatalog + now func() time.Time userOrigin string logger *slog.Logger } @@ -170,6 +182,10 @@ func NewEngine(opts Options) (*Engine, error) { if logger == nil { logger = slog.Default() } + now := opts.Now + if now == nil { + now = time.Now + } // The vendored lexicon catalog validates native postv2 input strictly: WE // sign the acceptance, so input that does not validate is fail-closed. Loaded // once at construction — a broken vendored file fails startup, not admission. @@ -189,6 +205,7 @@ func NewEngine(opts Options) (*Engine, error) { apActors: opts.APActors, maxPerCommunity: opts.MaxPerAuthorPerCommunity, catalog: catalog, + now: now, userOrigin: opts.UserOrigin, logger: logger, }, nil @@ -216,27 +233,39 @@ func (e *Engine) AdmitPost(ctx context.Context, did string, commit *consume.Comm } // The post's current binding: whether it was already accepted (so a now-failing - // re-admission is a REMOVAL, not a fresh rejection) and which community it is - // bound to (so a community-moving edit is discarded whole). The engine writes + // re-admission is a REMOVAL, not a fresh rejection). The engine writes // outbound_objects ONLY on accept, so a row here means the post federated. prior, priorBound, err := e.priorBinding(ctx, postURI) if err != nil { return err } + // The community this post is already bound to, read from EITHER surviving + // state — outbound_objects (accepted posts) OR the admissions ledger (a + // rejected post has a ledger row but no outbound row). A post is bound to one + // community forever; this is what makes community-immutability enforceable + // even for a post that was only ever rejected. + boundCommunity, err := e.boundCommunityOf(ctx, postURI, prior, priorBound) + if err != nil { + return err + } // Decide admission. The order is deliberate (fail closed first, cheap policy - // last): lexicon-validate → community-immutable → opt-out → paused → title → - // rate cap. A discard means the whole event is dropped (nothing written to - // either community); a non-empty code is a rejection/removal cause. - code, discard, err := e.decide(ctx, did, commit, communityDID, prior, priorBound) + // last): lexicon-validate → community-immutable → community-followed → opt-out + // → paused → title → rate cap. A discard means the whole event is dropped + // (nothing written to either community); a non-empty code is a rejection/ + // removal cause. + code, discard, err := e.decide(ctx, did, commit, communityDID, boundCommunity) if err != nil { return err } if discard { - e.logger.Debug("discarding community-moving edit", + // A community-moving edit: write NOTHING to the target community, and + // annotate the ORIGINAL community's ledger row (status unchanged) so the + // admin surface shows the attempted move instead of a silent drop. + e.logger.Info("discarding community-moving edit", slog.String("did", did), slog.String("post", postURI), - slog.String("bound_community", prior.CommunityDID), slog.String("event_community", communityDID)) - return nil + slog.String("bound_community", boundCommunity), slog.String("event_community", communityDID)) + return e.annotateCommunityImmutable(ctx, postURI) } if code != "" { @@ -259,38 +288,64 @@ func (e *Engine) AdmitPost(ctx context.Context, did string, commit *consume.Comm // The record body + context this decision was made against, so a later // force re-admit can re-run admission from stored state (a rejection // writes no outbound_objects, and the author's PDS is not local). - EvaluatedSnapshot: evaluatedSnapshot(commit), + EvaluatedSnapshot: e.evaluatedSnapshot(commit), }) } return e.accept(ctx, did, communityDID, postURI, commit) } -// operationDelete is the Jetstream commit operation for a record deletion. -const operationDelete = "delete" +// The AP op strings the deterministic activity id and the Page translation key +// on. operationDelete is also the Jetstream commit operation for a deletion; the +// others are DERIVED from outbound state, not read off the commit (see accept). +const ( + operationCreate = "create" + operationUpdate = "update" + operationDelete = "delete" +) // decide runs the admission checks in order and returns the rejection/removal // code ("" = admit), or discard=true when the event must be dropped whole (a // community-moving edit). The order fails closed first: garbage input never // reaches a policy check, and a hijack (community move) is refused before the // author's own preferences are consulted. -func (e *Engine) decide(ctx context.Context, did string, commit *consume.CommitEvent, communityDID string, - prior *store.OutboundObject, priorBound bool) (code string, discard bool, err error) { +func (e *Engine) decide(ctx context.Context, did string, commit *consume.CommitEvent, + communityDID, boundCommunity string) (code string, discard bool, err error) { - // 1. Strict lexicon validation of the native input — fail closed. - if !e.lexiconValid(commit.Record) { + // 1. Strict lexicon validation of the native input, bound to the postv2 + // schema — fail closed. A marshal/unmarshal fault is an INTERNAL error + // (retryable), NOT a permanent lexicon-invalid verdict. + valid, err := e.lexiconValid(commit.Record) + if err != nil { + return "", false, err + } + if !valid { return DecisionLexiconInvalid, false, nil } // 2. Community immutability. The lexicon marks `community` immutable: an - // UPDATE that names a different community than the post was accepted into is - // a retarget, which means writing a NEW post — so the whole event is - // discarded, not partially applied. - if priorBound && prior.CommunityDID != communityDID { + // UPDATE that names a different community than the post was bound to is a + // retarget, which means writing a NEW post — so the whole event is discarded, + // not partially applied. + if boundCommunity != "" && boundCommunity != communityDID { return "", true, nil } - // 3. Opt-out (decision 11): content pushed outward is exactly what an + // 3. Community follow gate (SECURITY): a communities row's mere existence is + // NOT authority to sign an acceptance into it. We federate only communities + // we hold an ACCEPTED Follow to; none/pending/unfollowed reject. + community, err := e.communities.GetByDID(ctx, communityDID) + if err != nil { + if errors.IsNotFound(err) { + return DecisionCommunityNotFollowed, false, nil + } + return "", false, fmt.Errorf("accept: resolve community %s: %w", communityDID, err) + } + if community.FollowState != store.FollowStateAccepted { + return DecisionCommunityNotFollowed, false, nil + } + + // 4. Opt-out (decision 11): content pushed outward is exactly what an // opted-out author refused. federating, err := e.mayFederate(ctx, did) if err != nil { @@ -300,7 +355,7 @@ func (e *Engine) decide(ctx context.Context, did string, commit *consume.CommitE return DecisionOptedOut, false, nil } - // 4. Paused (#account, decision 19): delivery is halted while the identity is + // 5. Paused (#account, decision 19): delivery is halted while the identity is // deactivated/suspended/takendown/throttled, so a new post is not admitted. if e.apActors != nil { actor, err := e.apActors.GetByDID(ctx, did) @@ -316,17 +371,19 @@ func (e *Engine) decide(ctx context.Context, did string, commit *consume.CommitE } } - // 5. Title: required and within Lemmy's cap (postv2 title is OPTIONAL in the - // lexicon, so this is admission policy, not validation). + // 6. Title: required and within Lemmy's cap (postv2 title is OPTIONAL in the + // lexicon, so this is admission policy, not validation). The cap counts RUNES, + // not bytes — Lemmy's limit is on grapheme length, so a multibyte title well + // under 200 characters must not be rejected for being over 200 bytes. title, _ := commit.Record["title"].(string) if title == "" { return DecisionTitleRequired, false, nil } - if len(title) > lemmyTitleCap { + if utf8.RuneCountInString(title) > lemmyTitleCap { return DecisionTitleTooLong, false, nil } - // 6. Rate cap: one author must not flood a community Tidepool vouches for. + // 7. Rate cap: one author must not flood a community Tidepool vouches for. // Counts the author's currently-accepted posts in this community, excluding // this post so a repin never counts against itself. 0 means unlimited. if e.maxPerCommunity > 0 { @@ -364,6 +421,18 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, return err } + // The AP op is a function of what LEMMY already holds, not of the Jetstream + // commit operation: Lemmy has a live copy only if a non-tombstoned outbound + // row already exists. An UPDATE of a never-federated post is a Create{Page}; + // an edit after a Delete (removal / author-delete then restore) is a Create + // too, because the live copy was withdrawn. This is read BEFORE the upsert + // bumps the row. + priorRow, priorErr := e.objects.GetByATURI(ctx, postURI) + if priorErr != nil && !errors.IsNotFound(priorErr) { + return fmt.Errorf("accept: read outbound state for %s: %w", postURI, priorErr) + } + wasLive := priorErr == nil && !priorRow.IsTombstoned() + snapshot, err := json.Marshal(map[string]any{ "atUri": postURI, "cid": commit.CID, @@ -398,10 +467,18 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, if err != nil { return fmt.Errorf("accept: write outbound state for %s: %w", postURI, err) } + // Create unless Lemmy already holds a live copy AND this is a later + // activity (seq bumped past the initial 0). The seq is guarded on CID + // change (UpsertTx), so an unchanged redelivery keeps seq 0 and reuses the + // original Create id rather than minting an Update to a peer. + op := operationCreate + if wasLive && stored.LastActivitySeq > 0 { + op = operationUpdate + } intent := consume.PostIntent{ - Op: commit.Operation, + Op: op, ATURI: postURI, - ID: consume.ActivityID(e.userOrigin, postURI, commit.Operation, stored.LastActivitySeq), + ID: consume.ActivityID(e.userOrigin, postURI, op, stored.LastActivitySeq), CommunityAPID: community.APGroupID, Snapshot: snapshot, } @@ -418,12 +495,26 @@ func (e *Engine) accept(ctx context.Context, did, communityDID, postURI string, EvaluatedCID: commit.CID, AcceptanceRKey: rkey, AcceptedCID: commit.CID, - EvaluatedSnapshot: evaluatedSnapshot(commit), + EvaluatedSnapshot: e.evaluatedSnapshot(commit), }) } - if _, err := acceptrec.AcceptSubject(ctx, e.repos, communityDID, postURI, commit.CID, - publishedAtOf(commit.Record), sideEffect); err != nil { + _, 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 + } + if err != nil { return fmt.Errorf("accept: admit %s into %s: %w", postURI, communityDID, err) } return nil @@ -463,15 +554,17 @@ func (e *Engine) removeAccepted(ctx context.Context, did, communityDID, postURI Status: StatusRemoved, DecisionCode: code, EvaluatedCID: commit.CID, - EvaluatedSnapshot: evaluatedSnapshot(commit), + EvaluatedSnapshot: e.evaluatedSnapshot(commit), }) } // The removal pins the version that was accepted when it was removed (audit // metadata); the code is the open-set admission-revoked, not a moderation - // reason — no moderator acted. + // reason — no moderator acted. createdAt is the DECISION time (the engine's + // clock), NOT the post's publication time — a removal is a claim about when we + // decided, so an old post removed today is dated today. if _, err := acceptrec.Remove(ctx, e.repos, communityDID, postURI, prior.LastCID, - RemovalCodeAdmissionRevoked, "", publishedAtOf(commit.Record), sideEffect); err != nil { + RemovalCodeAdmissionRevoked, "", e.now(), sideEffect); err != nil { return fmt.Errorf("accept: remove %s from %s: %w", postURI, communityDID, err) } return nil @@ -535,23 +628,95 @@ func (e *Engine) priorBinding(ctx context.Context, postURI string) (*store.Outbo return stored, true, nil } -// lexiconValid reports whether the postv2 record passes strict validation against -// the vendored lexicon catalog. Invalid input is fail-closed (rejected), never -// signed. A record with no $type, or one whose $type has no schema, is invalid. -func (e *Engine) lexiconValid(record map[string]any) bool { - recordType, _ := record["$type"].(string) - if recordType == "" { - return false +// lexiconValid reports whether the record is a valid social.coves.community +// .postv2, validated against THAT schema specifically. SECURITY: the type is +// bound to postv2, not trusted from the record's self-declared $type — WE sign +// the acceptance, so a profile- or comment-shaped record carrying a full postv2 +// body must not be signed as a community post just because it is a valid instance +// of the type it claims. +// +// The bool is the schema VERDICT (false → a real lexicon-invalid rejection). The +// error is an INTERNAL fault (marshal/unmarshal) — retryable, and NOT a permanent +// lexicon-invalid decision, so the caller must not record a rejection for it. +func (e *Engine) lexiconValid(record map[string]any) (bool, error) { + if t, _ := record["$type"].(string); t != materialize.CollectionPostV2 { + e.logger.Debug("record rejected: $type is not a postv2", + slog.String("type", t), slog.String("want", materialize.CollectionPostV2)) + return false, nil } raw, err := json.Marshal(record) if err != nil { - return false + return false, fmt.Errorf("accept: marshal record for validation: %w", err) } data, err := atdata.UnmarshalJSON(raw) if err != nil { - return false + return false, fmt.Errorf("accept: decode record for validation: %w", err) + } + if verr := lexicon.ValidateRecord(e.catalog, data, materialize.CollectionPostV2, lexicon.ValidateFlags(0)); verr != nil { + // A real schema verdict, logged with the offending field so an operator + // can tell a genuine bad record from a validator disagreement. + e.logger.Debug("record failed postv2 lexicon validation", + slog.String("detail", verr.Error())) + return false, nil + } + return true, nil +} + +// boundCommunityOf returns the community a post is already bound to, from EITHER +// surviving state: the outbound_objects row (an accepted post) or, failing that, +// the admissions ledger (a rejected post keeps a ledger row but no outbound row). +// "" means the engine has never decided on this post. This is what makes +// one-community-per-post enforceable even for a post that was only ever rejected. +func (e *Engine) boundCommunityOf(ctx context.Context, postURI string, priorRow *store.OutboundObject, priorBound bool) (string, error) { + if priorBound { + return priorRow.CommunityDID, nil + } + adm, err := e.admissions.GetByPostURI(ctx, postURI) + if errors.IsNotFound(err) { + return "", nil + } + if err != nil { + return "", err + } + return adm.CommunityDID, nil +} + +// annotateCommunityImmutable records the attempted community move on the post's +// EXISTING ledger row (status unchanged, decision_code = community-immutable) so +// the admin surface shows it, instead of a silent drop. It writes NOTHING to the +// target community — the row it updates is the one already keyed to the post's +// bound community. +func (e *Engine) annotateCommunityImmutable(ctx context.Context, postURI string) error { + adm, err := e.admissions.GetByPostURI(ctx, postURI) + if errors.IsNotFound(err) { + return nil // nothing decided yet; nothing to annotate + } + if err != nil { + return err + } + adm.DecisionCode = DecisionCommunityImmutable + return e.admissions.Record(ctx, *adm) +} + +// evaluatedSnapshot serializes the postv2 record and the context a Readmit needs +// to rebuild the CommitEvent it re-runs admission against. It is stored on EVERY +// decision (accept, reject, remove). A marshal failure is LOGGED (not silently +// masqueraded as a legacy row) and yields nil — the store coalesces that to '{}', +// which makes a later readmit surface as unrecoverable rather than re-run wrong. +func (e *Engine) evaluatedSnapshot(commit *consume.CommitEvent) []byte { + b, err := json.Marshal(map[string]any{ + "record": commit.Record, + "cid": commit.CID, + "rev": commit.Rev, + "operation": commit.Operation, + "collection": commit.Collection, + }) + if err != nil { + e.logger.Warn("failed to marshal evaluated snapshot; readmit will be unrecoverable for this decision", + slog.String("rkey", commit.RKey), slog.String("error", err.Error())) + return nil } - return lexicon.ValidateRecord(e.catalog, data, recordType, lexicon.ValidateFlags(0)) == nil + return b } // mayFederate reports whether the author permits outbound federation. A missing @@ -642,7 +807,11 @@ func (e *Engine) Readmit(ctx context.Context, postATURI string) (*ReadmitResult, if err != nil { return nil, err } - code, discard, err := e.decide(ctx, did, commit, communityDID, prior, priorBound) + boundCommunity, err := e.boundCommunityOf(ctx, postATURI, prior, priorBound) + if err != nil { + return nil, err + } + code, discard, err := e.decide(ctx, did, commit, communityDID, boundCommunity) if err != nil { return nil, err } @@ -669,7 +838,7 @@ func (e *Engine) Readmit(ctx context.Context, postATURI string) (*ReadmitResult, Status: StatusRejected, DecisionCode: code, EvaluatedCID: commit.CID, - EvaluatedSnapshot: evaluatedSnapshot(commit), + EvaluatedSnapshot: e.evaluatedSnapshot(commit), }); err != nil { return nil, err } @@ -683,24 +852,6 @@ func (e *Engine) Readmit(ctx context.Context, postATURI string) (*ReadmitResult, return &ReadmitResult{PostURI: postATURI, Status: StatusAccepted, Enqueued: true}, nil } -// evaluatedSnapshot serializes the postv2 record and the context a Readmit needs -// to rebuild the CommitEvent it re-runs admission against. It is stored on EVERY -// decision (accept, reject, remove). A marshal failure yields nil, which the -// store coalesces to '{}' — an unrecoverable readmit, never a wrong one. -func evaluatedSnapshot(commit *consume.CommitEvent) []byte { - b, err := json.Marshal(map[string]any{ - "record": commit.Record, - "cid": commit.CID, - "rev": commit.Rev, - "operation": commit.Operation, - "collection": commit.Collection, - }) - if err != nil { - return nil - } - return b -} - // rebuildCommit reconstructs the CommitEvent (and its author DID) a Readmit // re-runs admission against, from the post at-uri and the stored evaluated // snapshot. ok=false means the snapshot did not survive (legacy '{}' or a diff --git a/internal/accept/outer_acceptance_test.go b/internal/accept/outer_acceptance_test.go index 006adf2..30a7d86 100644 --- a/internal/accept/outer_acceptance_test.go +++ b/internal/accept/outer_acceptance_test.go @@ -165,6 +165,15 @@ func TestEngineCrashInjectionIsAtomic(t *testing.T) { assert.Zero(t, countRows(t, conn, "outbound_activities"), "and no outbound activity: the whole acceptance+enqueue transaction rolled back") assert.Zero(t, countRows(t, conn, "outbound_deliveries")) + + // The side effect writes the outbound_objects row AND the admissions ledger + // row on the SAME acceptance tx, so a rollback must leave NEITHER. A non-tx + // write of either would leak past the rollback undetected — pin both absent. + assert.Zero(t, countRows(t, conn, "outbound_objects"), + "the post's outbound_objects row must roll back with the acceptance (written on the tx)") + assert.Zero(t, countRows(t, conn, "admissions"), + "the accepted admissions row must roll back too — no ledger row may claim the post was "+ + "accepted when the acceptance itself did not commit") } // TestEngineRecordsOptedOutRejection pins the opt-out check MOVED into the diff --git a/internal/accept/review_test.go b/internal/accept/review_test.go new file mode 100644 index 0000000..3bdc4a8 --- /dev/null +++ b/internal/accept/review_test.go @@ -0,0 +1,397 @@ +package accept + +import ( + "context" + "database/sql" + "net/http" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/consume" + "tidepool/internal/store" +) + +// Round 4: the /second-opinion review findings — 6 HIGH + importants. + +// seedCommunityState registers a bridged community in a given follow state. +func seedCommunityState(t *testing.T, conn *sql.DB, did, apid, name string, state store.FollowState) { + t.Helper() + _, err := store.NewCommunities(conn).UpsertCommunity(context.Background(), store.Community{ + APGroupID: apid, + DID: did, + PreferredUsername: name, + Instance: acCommunityHost, + FollowState: state, + }) + require.NoError(t, err, "seed community in state %s", state) +} + +func admissionRowsForPost(t *testing.T, conn *sql.DB, postURI string) int { + t.Helper() + var n int + require.NoError(t, conn.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM admissions WHERE post_uri = $1`, postURI).Scan(&n)) + return n +} + +// --------------------------------------------------------------------------- +// H1 — FollowState gate (SECURITY): admission requires an ACCEPTED Follow +// --------------------------------------------------------------------------- + +func TestH1_AdmissionRequiresAcceptedFollowState(t *testing.T) { + for _, state := range []store.FollowState{store.FollowStateNone, store.FollowStatePending} { + t.Run(string(state), func(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedCommunityState(t, conn, acCommunityDID, acCommunityAPID, acCommunityName, state) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + + status, code := admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, StatusRejected, status, + "a post to a community Tidepool has no ACCEPTED Follow to must be rejected — a "+ + "communities row's mere existence is not authority to sign an acceptance into it") + assert.Equal(t, DecisionCommunityNotFollowed, code) + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.False(t, ok, "no acceptance is signed for an unfollowed community") + assert.Zero(t, countRows(t, conn, "outbound_activities"), "and nothing is enqueued") + }) + } + + t.Run("accepted control", func(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedCommunityState(t, conn, acCommunityDID, acCommunityAPID, acCommunityName, store.FollowStateAccepted) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + + status, _ := admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, StatusAccepted, status, "an ACCEPTED-follow community admits normally") + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.True(t, ok) + }) +} + +// --------------------------------------------------------------------------- +// H2 — lexicon $type binding (SECURITY): validate against postv2 specifically +// --------------------------------------------------------------------------- + +func TestH2_RecordMustBeAPostV2NotAnyValidLexicon(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + // A record whose $type is a DIFFERENT known lexicon (actor.profile has no + // required fields, so this validates against ITS schema) but carries a full + // postv2 body. The engine must validate against social.coves.community.postv2 + // SPECIFICALLY — trusting the record's self-declared $type would let a + // profile-shaped (or comment-shaped) record be signed as a community post. + confused := pv2Record(func(r map[string]any) { r["$type"] = "social.coves.actor.profile" }) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, confused))) + + status, code := admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, StatusRejected, status, + "a record that isn't a postv2 must be rejected even if it is a valid instance of its "+ + "own declared type — WE sign the acceptance, so the type is bound to the postv2 schema") + assert.Equal(t, DecisionLexiconInvalid, code) + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.False(t, ok, "no acceptance is signed over a non-postv2 record") +} + +// --------------------------------------------------------------------------- +// H3 — community immutability for REJECTIONS (reads the ADMISSIONS ledger) +// --------------------------------------------------------------------------- + +func TestH3_RejectedPostCannotBeMovedToAnotherCommunity(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) // community A (accepted follow) + seedBridgedCommunityB(t, conn) // community B (accepted follow) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + // Opted out → the create is REJECTED in community A. A rejection writes an + // admissions row (community_did=A) but NO outbound_objects row. + _, err := store.NewFederationPrefs(conn).Upsert(ctx, store.FederationPref{ + DID: acAuthorDID, Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + require.Equal(t, 1, admissionRowCount(t, conn, acCommunityDID, acPostURI), + "precondition: rejected in A") + require.Zero(t, countRows(t, conn, "outbound_objects"), + "precondition: a rejection writes no outbound_objects — so immutability can't read it") + + // The same post at-uri is now UPDATED targeting community B. The engine must + // read the prior community from the ADMISSIONS LEDGER (the only surviving + // state for a rejected post) and DISCARD the move. + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("update", acPostRKey, acRevUpdate, acPostCID2, acPostTimeUS+1, + pv2Record(func(r map[string]any) { r["community"] = acCommunityB_DID })))) + + assert.Zero(t, admissionRowCount(t, conn, acCommunityB_DID, acPostURI), + "a community-moving edit must write NOTHING to the target community, even for a "+ + "post that was only ever rejected") + assert.Equal(t, 1, admissionRowCount(t, conn, acCommunityDID, acPostURI), + "the original community's admission row is untouched") + assert.Equal(t, 1, admissionRowsForPost(t, conn, acPostURI), + "one post_uri must never hold admissions rows under two communities") +} + +// --------------------------------------------------------------------------- +// H4 — seq-bump idempotency: unchanged content must not mint a second activity +// --------------------------------------------------------------------------- + +func TestH4_ReadmitOfUnchangedAcceptedPostDoesNotMintASecondActivity(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + engine := engineWith(t, conn, repos, enq) + dispatcher := wireDispatcher(t, conn, engine, enq) + + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + require.Equal(t, 1, countRows(t, conn, "outbound_activities")) + originalID := consume.ActivityID(acUserOrigin, acPostURI, "create", 0) + + // Readmit the already-accepted, UNCHANGED post (a redelivery-shaped re-run). + _, err := engine.Readmit(ctx, acPostURI) + require.NoError(t, err) + + assert.Equal(t, 1, countRows(t, conn, "outbound_activities"), + "re-running an unchanged accepted post must NOT mint a second Create{Page}: the "+ + "outbound_objects seq must be guarded on content change (like TombstoneTx guards on "+ + "tombstoned_at), or every redelivery/readmit invents a new activity id") + var seq int + require.NoError(t, conn.QueryRowContext(ctx, + `SELECT last_activity_seq FROM outbound_objects WHERE at_uri = $1`, acPostURI).Scan(&seq)) + assert.Equal(t, 0, seq, "the activity seq must not bump when the content is unchanged") + + var kept string + require.NoError(t, conn.QueryRowContext(ctx, + `SELECT activity_id FROM outbound_activities`).Scan(&kept)) + assert.Equal(t, originalID, kept, "the one activity keeps its ORIGINAL id") +} + +// --------------------------------------------------------------------------- +// H5 — a corrective edit after an admission-revoked removal AUTO-RESTORES +// (RULING: propose auto-restore; the admission-revocation was OUR decision) +// --------------------------------------------------------------------------- + +func TestH5_CorrectiveEditAfterAdmissionRemovalAutoRestores(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + // Accepted, then edited titleless → REMOVED (admission-revoked removal stands). + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("update", acPostRKey, acRevUpdate, acPostCID2, acPostTimeUS+1, + pv2Record(func(r map[string]any) { delete(r, "title") })))) + _, removed := removalStandsAt(t, repos, acCommunityDID, acPostURI) + require.True(t, removed, "precondition: an admission-revoked removal stands") + + // A LATER corrective edit that now PASSES admission (title restored). It must + // NOT error out of AdmitPost (that would redrive forever against the standing + // removal); the admission-revocation was the bridge's own decision, so a + // corrective edit auto-restores. + err := dispatcher.HandleEvent(ctx, + postEvent("update", acPostRKey, "3lzpostrev004", acPostCID, acPostTimeUS+2, pv2Record())) + require.NoError(t, err, + "a corrective edit must be a DECIDED outcome, not a transient error that redrives forever "+ + "against ErrRemovalStands") + + _, stillRemoved := removalStandsAt(t, repos, acCommunityDID, acPostURI) + assert.False(t, stillRemoved, "the removal is deleted on auto-restore") + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.True(t, ok, "a fresh acceptance is written") + status, _ := admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, StatusAccepted, status, "the ledger is back to accepted") + + // Lemmy received a Delete{Page} at removal, so the restore re-adds it with a + // Create{Page} (op derived from state — the live copy is gone), not Update. + assert.Equal(t, 2, activityKindCount(t, conn, "Create"), + "the restore enqueues a Create{Page} (Lemmy has no live copy after the Delete)") + assert.Equal(t, 1, activityKindCount(t, conn, "Delete")) +} + +// --------------------------------------------------------------------------- +// H6 — op derived from outbound STATE, not the Jetstream commit operation +// --------------------------------------------------------------------------- + +func TestH6_UpdateEventForNeverAcceptedPostEmitsCreate(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + // An UPDATE event for a post Lemmy has NEVER seen (no outbound state). Lemmy + // got no Create, so this must federate as Create{Page}, not Update{Page} — + // the op is a function of what Lemmy already holds, not of the commit operation. + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("update", acPostRKey, acRevUpdate, acPostCID, acPostTimeUS, pv2Record()))) + + assert.Equal(t, 1, activityKindCount(t, conn, "Create"), + "an update of a never-federated post is a Create{Page}: op comes from outbound-state "+ + "presence, not commit.Operation") + assert.Zero(t, activityKindCount(t, conn, "Update"), + "Lemmy never received a Create, so an Update would reference an object it does not have") + wantCreateID := consume.ActivityID(acUserOrigin, acPostURI, "create", 0) + var id string + require.NoError(t, conn.QueryRowContext(ctx, + `SELECT activity_id FROM outbound_activities`).Scan(&id)) + assert.Equal(t, wantCreateID, id, "the activity id uses op=create") +} + +func TestH6_UpdateOfAcceptedPostEmitsUpdate(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, pv2Record()))) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("update", acPostRKey, acRevUpdate, acPostCID2, acPostTimeUS+1, pv2Record()))) + + assert.Equal(t, 1, activityKindCount(t, conn, "Update"), + "an edit of an accepted post (Lemmy has a live copy) federates as Update{Page}") +} + +// --------------------------------------------------------------------------- +// Important — title cap counts RUNES, not bytes +// --------------------------------------------------------------------------- + +func TestImportant_TitleCapCountsRunesNotBytes(t *testing.T) { + t.Run("150 multibyte runes accepted", func(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, repos, enq), enq) + + // 150 CJK runes = 450 bytes: over the byte cap but well under the 200-rune + // (grapheme) cap Lemmy actually enforces. + title := strings.Repeat("あ", 150) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, + pv2Record(func(r map[string]any) { r["title"] = title })))) + + status, code := admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, StatusAccepted, status, + "a 150-rune multibyte title (450 bytes) must be ACCEPTED: the cap counts runes, not bytes") + assert.Empty(t, code) + _, ok := acceptanceSubjectCID(t, repos, acCommunityDID, acPostURI) + assert.True(t, ok) + }) + + t.Run("201 runes rejected", func(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + enq := realEnqueuer(t, conn) + dispatcher := wireDispatcher(t, conn, engineWith(t, conn, newRepos(t, conn), enq), enq) + + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, + pv2Record(func(r map[string]any) { r["title"] = strings.Repeat("x", lemmyTitleCap+1) })))) + + _, code := admissionOf(t, conn, acCommunityDID, acPostURI) + assert.Equal(t, DecisionTitleTooLong, code, "201 runes exceeds the cap") + }) +} + +// --------------------------------------------------------------------------- +// Important — a removal's createdAt is the DECISION time, not the post's +// --------------------------------------------------------------------------- + +func TestImportant_RemovalCreatedAtIsDecisionTimeNotPublication(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + seedBridgedCommunity(t, conn) + repos := newRepos(t, conn) + enq := realEnqueuer(t, conn) + + decisionTime := time.Date(2026, 8, 13, 12, 0, 0, 0, time.UTC) + engine := engineWith(t, conn, repos, enq, func(o *Options) { o.Now = func() time.Time { return decisionTime } }) + dispatcher := wireDispatcher(t, conn, engine, enq) + + // An OLD post (published in 2020), accepted, then edited titleless → removed. + oldPublished := "2020-01-01T00:00:00.000Z" + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("create", acPostRKey, acRevCreate, acPostCID, acPostTimeUS, + pv2Record(func(r map[string]any) { r["createdAt"] = oldPublished })))) + require.NoError(t, dispatcher.HandleEvent(ctx, + postEvent("update", acPostRKey, acRevUpdate, acPostCID2, acPostTimeUS+1, + pv2Record(func(r map[string]any) { r["createdAt"] = oldPublished; delete(r, "title") })))) + + removal, ok := removalStandsAt(t, repos, acCommunityDID, acPostURI) + require.True(t, ok) + assert.Equal(t, "2026-08-13T12:00:00.000Z", removal["createdAt"], + "a removal's createdAt is WHEN THE ENGINE DECIDED (injected now), not the post's 2020 "+ + "publication time — a removal is a claim about the decision, not the content") +} + +// --------------------------------------------------------------------------- +// Admin — readmit error mapping (regression pins for the existing handler) +// --------------------------------------------------------------------------- + +func TestAdminReadmit_UnknownPostIs404(t *testing.T) { + conn := acceptanceDB(t) + engine := engineWith(t, conn, newRepos(t, conn), realEnqueuer(t, conn)) + router := newAdminRouter(t, conn, engine) + + rec := adminRequest(t, router, http.MethodPost, "/admin/admissions/readmit", adminToken, + `{"post":"at://did:plc:nobody/social.coves.community.postv2/nope"}`) + assert.Equal(t, http.StatusNotFound, rec.Code, "readmit of a post with no admission row is 404") +} + +func TestAdminReadmit_LegacyEmptySnapshotIs422(t *testing.T) { + conn := acceptanceDB(t) + ctx := context.Background() + // A ledger row with NO stored snapshot (a legacy/pre-migration-022 decision). + require.NoError(t, NewAdmissions(conn).Record(ctx, Admission{ + AuthorDID: acAuthorDID, + CommunityDID: acCommunityDID, + PostURI: acPostURI, + Status: StatusRejected, + DecisionCode: DecisionOptedOut, + // EvaluatedSnapshot left nil → stored as '{}'. + })) + engine := engineWith(t, conn, newRepos(t, conn), realEnqueuer(t, conn)) + router := newAdminRouter(t, conn, engine) + + rec := adminRequest(t, router, http.MethodPost, "/admin/admissions/readmit", adminToken, + `{"post":"`+acPostURI+`"}`) + assert.Equal(t, http.StatusUnprocessableEntity, rec.Code, + "a readmit whose record snapshot did not survive is a 422, never a silent no-op") +} diff --git a/internal/acceptrec/acceptrec.go b/internal/acceptrec/acceptrec.go index b75ff7a..5c35355 100644 --- a/internal/acceptrec/acceptrec.go +++ b/internal/acceptrec/acceptrec.go @@ -211,11 +211,21 @@ func DeleteAcceptance(ctx context.Context, repos RepoManager, communityDID, subj // inside the commit. func Remove(ctx context.Context, repos RepoManager, communityDID, subjectURI, subjectCID, code, reason string, at time.Time, sideEffect repo.TxSideEffect) (*repo.CommitResult, error) { rkey := SubjectRKey(subjectURI) + // createdAt is the DECISION time (`at`, the engine's clock), but a standing + // removal carries its createdAt FORWARD so a redelivery re-puts byte-identical + // bytes and the repo layer's NoOp path absorbs it — only the first write + // stamps the clock. + createdAt := recordDatetime(at) + if existing, _, err := repos.GetRecord(ctx, communityDID, CollectionRemoval, rkey); err == nil { + if when, ok := existing["createdAt"].(string); ok && when != "" { + createdAt = when + } + } removal := map[string]any{ "$type": CollectionRemoval, "subject": strongRef(subjectURI, subjectCID), "code": code, - "createdAt": recordDatetime(at), + "createdAt": createdAt, } // Omitted rather than written blank: an empty reason renders in a moderation // log as a blank explanation instead of as none given. @@ -234,3 +244,31 @@ func Remove(ctx context.Context, repos RepoManager, communityDID, subjectURI, su } return res, nil } + +// Restore is the inverse of Remove: it deletes the standing removal and writes a +// fresh acceptance at the shared rkey in ONE commit, running sideEffect (the +// re-delivery enqueue) inside it. It is what a corrective edit routes through +// after an admission-revoked removal — the removal was OUR decision, so a post +// that now passes admission is reinstated rather than left withdrawn. Unlike +// AcceptSubject, it does NOT trip the removal guard: deleting the removal is the +// point. createdAt is derived from publishedAt (a fresh acceptance stands for the +// current version), so a redelivery re-puts byte-identical bytes. +func Restore(ctx context.Context, repos RepoManager, communityDID, subjectURI, subjectCID string, publishedAt time.Time, sideEffect repo.TxSideEffect) (*repo.CommitResult, error) { + rkey := SubjectRKey(subjectURI) + acceptance := map[string]any{ + "$type": CollectionAcceptance, + "subject": strongRef(subjectURI, subjectCID), + "createdAt": recordDatetime(publishedAt), + } + // One commit: the removal is withdrawn and the acceptance written together, + // so the firehose never shows a window where the post is neither removed nor + // accepted. + res, err := repos.ApplyOpsTx(ctx, communityDID, []repo.RecordOp{ + {Action: repo.OpActionDelete, Collection: CollectionRemoval, RKey: rkey}, + {Action: repo.OpActionUpdate, Collection: CollectionAcceptance, RKey: rkey, Record: acceptance}, + }, sideEffect) + if err != nil { + return nil, fmt.Errorf("acceptrec: restore %s into %s: %w", subjectURI, communityDID, err) + } + return res, nil +} diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index 6f4329a..3f9b0ac 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -481,9 +481,8 @@ func (d *Dispatcher) handlePostV2(ctx context.Context, _ *sql.Tx, did string, co // (decision_code opted-out) in the admissions ledger rather than the // consumer dropping it silently at debug. post.getStatus and the admin // surface both need the "why", and only the engine writes it. The - // consumer keeps only the pre-gate an opted-out author's post still fails - // for a DIFFERENT reason: it must name a bridged community for the engine - // to have a repo to reject it INTO. + // consumer keeps just one pre-gate check: a post must name a bridged + // community for the engine to have a repo to reject it INTO. communityDID := stringField(commit.Record, "community") if communityDID == "" { // The lexicon REQUIRES community. A post without one is malformed diff --git a/internal/repo/applyops_tx_test.go b/internal/repo/applyops_tx_test.go index b5c0337..22cf22e 100644 --- a/internal/repo/applyops_tx_test.go +++ b/internal/repo/applyops_tx_test.go @@ -176,3 +176,33 @@ func TestApplyOpsTx_SuccessRunsSideEffectInSameCommit(t *testing.T) { _, _, err = manager.GetRecord(ctx, testDID, testCollection, testRKey(2)) assert.NoError(t, err, "the written record exists") } + +// TestApplyOpsTx_GenesisDeleteRunsSideEffect: an all-delete batch against a +// community that has NO repo yet (state == nil, KeyUseDelete) currently returns +// early NoOp WITHOUT running the side effect. Per the at-least-once side-effect +// contract (the enqueue must fire even when the repo commit is inert), the side +// effect MUST still run and commit — a DeleteAcceptance-shaped retraction whose +// community repo was never created must still enqueue its Delete{Page}. +// +// RULING (flagged): the side effect runs here too, for the SAME reason the +// len(emitted)==0 branch runs it. If the coordinator rules this an intentional +// no-side-effect exit instead, this pin is the place to invert. +func TestApplyOpsTx_GenesisDeleteRunsSideEffect(t *testing.T) { + manager, database, _ := testManager(t) + applyOpsTxMarker(t, database) + ctx := context.Background() + + // testDID has no repo_state row (no PutRecord ran), so state == nil and the + // all-delete batch takes the genesis-delete branch. + res, err := manager.ApplyOpsTx(ctx, testDID, []RecordOp{ + {Action: OpActionDelete, Collection: testCollection, RKey: testRKey(1)}, + {Action: OpActionDelete, Collection: testOtherCollection, RKey: testRKey(2)}, + }, writeMarker(ctx, "genesis-delete-enqueue", nil)) + require.NoError(t, err) + require.NotNil(t, res) + assert.True(t, res.NoOp, "no records exist, so the batch commits no records") + + assert.Equal(t, 1, markerCount(t, database), + "the side effect must run on the genesis-delete branch too: the at-least-once outbound "+ + "enqueue has to fire even when the community repo was never created") +} diff --git a/internal/repo/repo.go b/internal/repo/repo.go index 66dfd45..f57537c 100644 --- a/internal/repo/repo.go +++ b/internal/repo/repo.go @@ -361,7 +361,20 @@ func (m *Manager) ApplyOpsTx(ctx context.Context, did string, ops []RecordOp, si // Nothing exists and nothing is being written: every op is a // delete against a repo with no records, so the whole batch is // inert. Genesis is reserved for commits that put something. - return &CommitResult{NoOp: true}, nil + res := &CommitResult{NoOp: true} + if sideEffect != nil { + // The at-least-once side effect (the outbound enqueue) still runs + // and must be made durable even though no repo commit happened — + // a retraction whose community repo was never created must still + // enqueue its Delete{Page}. Mirrors the len(emitted)==0 branch. + if err := sideEffect(ctx, tx, res); err != nil { + return nil, err + } + if err := tx.Commit(); err != nil { + return nil, fmt.Errorf("repo: commit genesis-delete side effect for %s: %w", did, err) + } + } + return res, nil } empty := mst.NewEmptyTree() tree = &empty diff --git a/internal/store/outbound_objects.go b/internal/store/outbound_objects.go index 6f142f2..37030e3 100644 --- a/internal/store/outbound_objects.go +++ b/internal/store/outbound_objects.go @@ -65,7 +65,19 @@ func (r *postgresOutboundObjects) upsert(ctx context.Context, q execer, object O community_ap_id = EXCLUDED.community_ap_id, translated_snapshot = EXCLUDED.translated_snapshot, depth = EXCLUDED.depth, - last_activity_seq = outbound_objects.last_activity_seq + 1, + -- The seq bumps ONLY when the commit actually changed — the record CID + -- OR the commit rev — mirroring how TombstoneTx guards its bump on + -- tombstoned_at IS NULL. A redelivery or a force re-admit of UNCHANGED + -- content re-runs the outbound side effect against byte-identical state + -- (same cid AND same rev); it must reuse the same activity id, or every + -- replay invents a new one and re-delivers the same object. A real edit + -- (new cid, and always a new rev) still bumps — including an edit that + -- happens to reserialize to the same cid but under a fresh rev. + last_activity_seq = CASE + WHEN outbound_objects.last_cid IS DISTINCT FROM EXCLUDED.last_cid + OR outbound_objects.last_rev IS DISTINCT FROM EXCLUDED.last_rev + THEN outbound_objects.last_activity_seq + 1 + ELSE outbound_objects.last_activity_seq END, updated_at = now() RETURNING` + outboundObjectColumns