diff --git a/internal/consume/comments.go b/internal/consume/comments.go index bb06d20..2e7a4aa 100644 --- a/internal/consume/comments.go +++ b/internal/consume/comments.go @@ -37,6 +37,9 @@ func (d *Dispatcher) handleComment(ctx context.Context, tx *sql.Tx, did string, case operationDelete: return d.applyCommentDelete(ctx, tx, did, commit) default: + // Unreachable: validateCommitEnvelope rejects any other operation as + // permanent before the gate. A claimed skip, since an operation this + // build cannot name is not something a replay improves. d.logger.Debug("unknown comment operation", slog.String("operation", commit.Operation), slog.String("did", did)) return nil @@ -68,6 +71,10 @@ func (d *Dispatcher) applyCommentWrite(ctx context.Context, tx *sql.Tx, did stri if !federating { // The residual split-thread case, explicitly chosen (decision 11): the // comment stays on the atproto side and the Lemmy side never sees it. + // + // PERMANENT, and it CLAIMS the gate: the author's "no" is not a fact + // that changes underneath this frame, and a later re-enable federates + // what they write next rather than what they wrote while opted out. d.logger.Debug("skipping comment from an opted-out author", slog.String("did", did), slog.String("rkey", commit.RKey), slog.String("operation", commit.Operation)) @@ -84,10 +91,16 @@ func (d *Dispatcher) applyCommentWrite(ctx context.Context, tx *sql.Tx, did stri // with no prior state are edits to a comment that never federated. // Dead-lettering either would bury the queue in events working exactly // as intended. - d.logger.Debug("skipping comment with no federated thread", - slog.String("did", did), slog.String("rkey", commit.RKey), - slog.String("operation", commit.Operation)) - return nil + // + // UNCLAIMED. A reply can reach this consumer before the thing it replies + // to has been materialized — the two arrive milliseconds apart and + // nothing orders them — and a claimed gate row would swallow the + // reconnect rewind that recovers it, so the comment would silently never + // federate. The update case rides the same ruling for the same reason: + // an edit resolves through the state its CREATE wrote, so it becomes + // resolvable exactly when the create is recovered. + return fmt.Errorf("%w: comment %s belongs to no thread this bridge federates (yet)", + errSkipUnclaimed, atURI) } if thread.Depth > maxCommentDepth { @@ -280,6 +293,11 @@ func (d *Dispatcher) applyCommentDelete(ctx context.Context, tx *sql.Tx, did str if errors.IsNotFound(err) { // A comment this bridge never federated. There is nothing to withdraw, // and most native comment deletes are exactly this. + // + // PERMANENT, and the claim is LOAD-BEARING: the gate row a delete leaves + // IS the tombstone that rejects a stale create for the same URI, so + // releasing it would let a replay past the delete resurrect content the + // user removed. d.logger.Debug("skipping delete for a comment with no outbound state", slog.String("did", did), slog.String("rkey", commit.RKey)) return nil diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index ad6f1f5..1f37dcb 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -517,6 +517,18 @@ func validateCommitEnvelope(did string, commit *CommitEvent) error { // here is also what makes the lexicon's community-immutability rule // enforceable in ONE place: there is no second copy of the answer to disagree // with the engine's. +// +// THE GATE TRANSACTION IS DELIBERATELY NOT PASSED ON, and the underscore says +// so rather than hiding it. AdmitPost takes no transaction and opens its own — +// the acceptance write into the community repo, the ledger row and the outbound +// enqueue are one commit that only the engine can compose — so handing it this +// one would break the rule rev_gate.go's DEADLOCK NOTE states: a handler must +// not write on the gate tx and then call something that opens a second +// transaction touching the same rows. The consequence is stated plainly: if the +// gate later fails to commit, the engine's work stands under an unadvanced gate +// and the replay re-enters AdmitPost. That is survivable only because the +// engine decides from STORED state (its prior binding, its ledger row), so a +// re-admission of the same rev re-reaches the same conclusion. func (d *Dispatcher) handlePostV2(ctx context.Context, _ *sql.Tx, did string, commit *CommitEvent) error { // The nil-engine skip happens before the gate (see handleCommit), so a // non-nil engine is guaranteed here. @@ -551,9 +563,14 @@ func (d *Dispatcher) handlePostV2(ctx context.Context, _ *sql.Tx, did string, co // Most Coves posts are exactly this. Admitting one would write an // acceptance record into a community repo that has no business // existing. - d.logger.Debug("skipping postv2 for a non-bridged community", - slog.String("did", did), slog.String("community", communityDID)) - return nil + // + // UNCLAIMED, for the same reason the nil-engine skip happens BEFORE + // the gate a few lines up in handleCommit: whether this community is + // bridged is an operator's decision that can be made tomorrow, and a + // gate row claimed today would turn the replay that should admit the + // post into a silent no-op, dropping it forever. + return fmt.Errorf("%w: postv2 %s is for community %s, which this bridge does not federate (yet)", + errSkipUnclaimed, commitRecordURI(did, commit), strconv.Quote(communityDID)) } } diff --git a/internal/consume/fault_test.go b/internal/consume/fault_test.go index 152ea5e..f1695ea 100644 --- a/internal/consume/fault_test.go +++ b/internal/consume/fault_test.go @@ -43,6 +43,41 @@ func (o *failingOutboundObjects) GetByATURI(ctx context.Context, atURI string) ( return o.OutboundObjects.GetByATURI(ctx, atURI) } +// failingOutboundDeliveries delegates to a real store but fails the opt-out's +// cancellation. That call is the LAST thing the soft opt-out does on the +// rev-gate transaction before it commits, so failing it is the realistic shape +// of "the preference was written, then the unit failed" — the window in which a +// preference written on its own connection would survive an unadvanced gate. +type failingOutboundDeliveries struct { + store.OutboundDeliveries + err error +} + +func (d *failingOutboundDeliveries) CancelOutwardForActorTx(_ context.Context, _ *sql.Tx, _ string) (int64, error) { + return 0, d.err +} + +// prefProbingDeleter stands in for the destructive seam and answers the one +// question that seam's ORDERING is about: by the time peers are asked to delete +// this user's content, is the user's intent already durable? It reads +// federation_prefs on its OWN connection, exactly as outbound.Purger does. +type prefProbingDeleter struct { + db *sql.DB + called bool + sawPreference bool +} + +func (d *prefProbingDeleter) DeleteRemoteContent(ctx context.Context, did string) error { + d.called = true + var count int + if err := d.db.QueryRowContext(ctx, + `SELECT COUNT(*) FROM federation_prefs WHERE did = $1`, did).Scan(&count); err != nil { + return err + } + d.sawPreference = count > 0 + return nil +} + // failingCommunities delegates to a real store but fails GetByDID. The vote // delete path resolves the target community through it (communityAPID), where // a swallowed DB error silently addresses the Undo to nobody. diff --git a/internal/consume/federation.go b/internal/consume/federation.go index d2ab9a8..11640dc 100644 --- a/internal/consume/federation.go +++ b/internal/consume/federation.go @@ -31,12 +31,16 @@ import ( // clears any stored deleteRemote: a stale destructive flag on a re-enabled // user is a loaded gun pointed at the task 17 tier. // -// tx is the REV-GATE's transaction, and the disable path writes on it: the -// actor mirror and the cancellation of that actor's queued deliveries are one -// decision, and they commit with the gate advance or not at all. Atomicity here -// is structural rather than defended — there is no window in which one landed -// and the other did not, and a failure leaves the gate un-advanced so the record -// replays and re-applies both. +// tx is the REV-GATE's transaction, and every state write this handler makes +// rides it: the preference row, the cancellation of that actor's queued +// deliveries, and the actor mirror are ONE decision that commits with the gate +// advance or not at all. A failure leaves the gate un-advanced and nothing +// committed, so the record replays and re-applies the whole decision. +// +// WITH ONE EXCEPTION, and it is the destructive tier — see applyOptOut. The +// preference has to be durable BEFORE peers are asked to delete anything, and +// the seam that asks them writes the very same row from its own transaction, so +// on that path the preference is committed ahead of the gate on purpose. func (d *Dispatcher) handleFederation(ctx context.Context, tx *sql.Tx, did string, commit *CommitEvent) error { if commit.Operation == operationDelete { // A delete commit carries no record body, which costs nothing here: @@ -60,19 +64,53 @@ func (d *Dispatcher) handleFederation(ctx context.Context, tx *sql.Tx, did strin // deleteRemote is optional and defaults to false: the soft tier. Nothing // destructive is ever INFERRED — only an explicit true escalates. deleteRemote, _ := boolField(commit.Record, "deleteRemote") + return d.applyOptOut(ctx, tx, did, deleteRemote) +} - // The preference is recorded BEFORE the destructive seam is reached, and - // deliberately NOT on the transaction below. Peers that honor a Delete - // cannot restore what they dropped, so the user's intent must be durable - // before anything is sent — and the gate transaction has not committed by - // the time the destructive seam runs. It is the authority besides: it - // answers for every DID, including the ones with no actor to mirror onto. - if _, err := d.prefs.Upsert(ctx, store.FederationPref{ +// applyOptOut records the opt-out and stops the user's outbound traffic. +// +// THE RESULTING ORDER, and which connection each step runs on, because that is +// the whole of this function: +// +// 1. the PREFERENCE. On the gate tx for the soft tier; on its own connection, +// committed immediately, for the destructive one (see below). +// 2. the CANCELLATION of everything already queued — on the gate tx. +// 3. the DESTRUCTIVE SEAM, if asked for. Outside any transaction of ours: it +// opens its own. +// 4. the ACTOR MIRROR — on the gate tx, and LAST, because the seam above +// updates that same row from its own transaction. +// +// WHY STEP 1 SPLITS. The seam in step 3 reads and writes federation_prefs from +// its own transaction: outbound.Purger marks this exact row purged there, and +// purged_at is the only durable record that peers were really asked to delete. +// Holding the row uncommitted on the gate tx across that call would do both +// halves of the damage at once — the purge would find no preference to mark, +// and its UPDATE of the row would block on a lock this handler cannot release +// without committing while it waits synchronously for the call to return. That +// is precisely the shape rev_gate.go's DEADLOCK NOTE forbids. +// +// The cost of the split is stated rather than hidden: on the destructive path a +// later failure leaves a committed preference under an unadvanced gate, so the +// record replays. That is safe because the opt-out is idempotent (the same +// upsert, the same deterministic activity ids, a delivery insert that returns +// the standing row) and because the preference is the SAFE half to have +// committed early — it says "stop", and a replay re-asserts it. +func (d *Dispatcher) applyOptOut(ctx context.Context, tx *sql.Tx, did string, deleteRemote bool) error { + // The preference is the AUTHORITY, not a mirror: it answers for every DID, + // including the ones with no actor to mirror onto. + pref := store.FederationPref{ DID: did, Enabled: false, DeleteRemote: deleteRemote, Source: store.FederationPrefSourceRecord, - }); err != nil { + } + var err error + if deleteRemote { + _, err = d.prefs.Upsert(ctx, pref) + } else { + _, err = d.prefs.UpsertTx(ctx, tx, pref) + } + if err != nil { return fmt.Errorf("record federation opt-out for %s: %w", did, err) } @@ -165,21 +203,38 @@ func (d *Dispatcher) deleteRemoteContent(ctx context.Context, did string) error // Delete cannot remove a purged preference — so this reads the outcome back // rather than deciding it, and says so once, loudly, where an operator can see // that a user tried to come back and could not. +// +// BOTH HALVES RIDE THE GATE TRANSACTION. No destructive seam is reachable from +// here, so nothing opens a second transaction against either row and the whole +// restore commits with the gate advance — which matters in this direction most +// of all: absence means default-on, so a cleared preference that outlived a +// failed event would silently re-enable federation for somebody who asked us to +// stop. func (d *Dispatcher) restoreDefaultFederation(ctx context.Context, tx *sql.Tx, did string) error { - if err := d.prefs.Delete(ctx, did); err != nil { + cleared, err := d.prefs.DeleteTx(ctx, tx, did) + if err != nil { return fmt.Errorf("clear federation preference for %s: %w", did, err) } if err := d.mirrorActorEnabled(ctx, tx, did, true); err != nil { return err } + if cleared { + return nil // an ordinary re-enable + } + // Nothing was cleared, so this is either a DID that never opted out or one + // whose withdrawal committed and cannot be undone. The read tells them + // apart, and it is safe on any connection precisely BECAUSE nothing was + // written: this transaction has changed nothing about the row it is asking + // about, so an outside snapshot and the transaction's own agree. + // // The read is the report. Nothing here can fail the event: the record was // applied exactly as far as it is allowed to go, and retrying would re-ask a // question whose answer is terminal. pref, err := d.prefs.Get(ctx, did) switch { case errors.IsNotFound(err): - return nil // cleared: an ordinary re-enable + return nil // there was nothing to clear case err != nil: return fmt.Errorf("read federation preference for %s: %w", did, err) case pref.PurgedAt != nil: diff --git a/internal/consume/federation_atomicity_test.go b/internal/consume/federation_atomicity_test.go new file mode 100644 index 0000000..f4e61ed --- /dev/null +++ b/internal/consume/federation_atomicity_test.go @@ -0,0 +1,114 @@ +package consume + +import ( + "context" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// Chunk 3 finding 4: gate atomicity reached only two of the four handlers. The +// federation handler wrote its PREFERENCE row on its own autocommit connection +// while the actor mirror and the delivery cancellation rode the rev-gate +// transaction — so a failure between them left the preference committed under +// an unadvanced gate, and the handler's own doc comment asserted the opposite. +// +// These pin the two halves separately, because the two paths have different +// rulings and the difference is the whole design (see handleFederation): +// +// - the SOFT opt-out and the re-enable write their state ON the gate tx; +// - the DESTRUCTIVE opt-out deliberately does not, because the seam it calls +// opens its own transaction against the same row. + +func TestFederationAtomicity_SoftOptOutPreferenceRollsBackWithTheGate(t *testing.T) { + database := dispatchTestDB(t) + ctx := context.Background() + + // The cancellation runs AFTER the preference write, on the gate tx. Failing + // it is the realistic shape of "everything up to here committed, then the + // unit failed". + deliveries := &failingOutboundDeliveries{ + OutboundDeliveries: store.NewOutboundDeliveries(database), + err: fmt.Errorf("statement timeout"), + } + fixture := newDispatchFixture(t, database, func(opts *Options) { opts.Deliveries = deliveries }) + + err := fixture.handle(t, federationFrameNoDeleteRemote(dispatchNativeDID, dispatchRev, false)) + require.Error(t, err, "a failed cancellation must fail the event so it replays") + + assert.Zero(t, countRows(t, database, "federation_prefs"), + "NO preference row may survive: it is the handler's durable state and it rides "+ + "the gate transaction, so a failure rolls it back with the gate advance") + assert.Zero(t, countRows(t, database, "jetstream_record_revs"), + "and the gate is unadvanced, so the replay re-enters the handler") + + // And the replay recovers the whole decision. + replay := newDispatchFixture(t, database) + require.NoError(t, replay.handle(t, + federationFrameNoDeleteRemote(dispatchNativeDID, dispatchRev, false))) + pref, err := store.NewFederationPrefs(database).Get(ctx, dispatchNativeDID) + require.NoError(t, err) + assert.False(t, pref.Enabled) +} + +func TestFederationAtomicity_ReEnableClearsThePreferenceOnTheGateTx(t *testing.T) { + database := dispatchTestDB(t) + ctx := context.Background() + prefs := store.NewFederationPrefs(database) + + _, err := prefs.Upsert(ctx, store.FederationPref{ + DID: dispatchNativeDID, + Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + + fixture := newDispatchFixture(t, database) + + // White-box on purpose: the re-enable path has no injectable seam AFTER the + // preference write, so the only way to observe which connection the delete + // ran on is to hand it a transaction and roll that transaction back. If the + // delete rode its own autocommit connection the row would be gone anyway. + tx, err := database.BeginTx(ctx, nil) + require.NoError(t, err) + require.NoError(t, fixture.dispatcher.restoreDefaultFederation(ctx, tx, dispatchNativeDID)) + require.NoError(t, tx.Rollback()) + + pref, err := prefs.Get(ctx, dispatchNativeDID) + require.NoError(t, err, + "the preference SURVIVES a rolled-back gate transaction: absence means "+ + "default-on, so a delete that outlived a failed event would silently "+ + "re-enable federation for a user who asked us to stop") + require.NotNil(t, pref) + assert.False(t, pref.Enabled) +} + +func TestFederationAtomicity_DestructiveOptOutCommitsThePreferenceBeforeTheSeam(t *testing.T) { + database := dispatchTestDB(t) + ctx := context.Background() + + // The destructive seam observes federation_prefs from its OWN transaction + // (outbound.Purger marks the row purged there). If the preference were held + // uncommitted on the gate tx, that read would miss it — and the purge's + // UPDATE of the same row would block on a lock the handler cannot release + // without committing, which is the deadlock rev_gate.go's DEADLOCK NOTE + // names. + probe := &prefProbingDeleter{db: database} + fixture := newDispatchFixture(t, database, func(opts *Options) { opts.RemoteDeleter = probe }) + + require.NoError(t, fixture.handle(t, + federationFrame(dispatchNativeDID, dispatchRev, "create", false, true))) + + require.True(t, probe.called, "the destructive seam ran") + assert.True(t, probe.sawPreference, + "the user's intent is DURABLE before anything irreversible is sent: peers that "+ + "honour a Delete cannot restore what they dropped, and the purge marks this "+ + "very row from its own transaction") + + pref, err := store.NewFederationPrefs(database).Get(ctx, dispatchNativeDID) + require.NoError(t, err) + assert.True(t, pref.DeleteRemote) +} diff --git a/internal/consume/metrics.go b/internal/consume/metrics.go index 4ca67e4..b2141b2 100644 --- a/internal/consume/metrics.go +++ b/internal/consume/metrics.go @@ -21,8 +21,31 @@ const ( MetricConnected = "tidepool_consumer_connected" MetricDialFailures = "tidepool_consumer_dial_failures" MetricDisconnectedSeconds = "tidepool_consumer_disconnected_seconds" + MetricUnclaimedSkips = "tidepool_consumer_unclaimed_skips" ) +// unclaimedSkips counts the handler skips that DELIBERATELY release the rev +// gate un-advanced (errSkipUnclaimed): a vote whose subject has not +// materialized, a comment whose thread has not, a post for a community nobody +// has bridged yet, a vote direction this build does not understand. +// +// It is the only trace such an event leaves. It is not a failure, so it never +// reaches the DLQ; it is not applied, so it writes no state; and the event is +// reported handled, so the cursor moves past it. A backlog of them is a real +// condition — a community that should have been bridged, a materializer that +// has stalled — and without this counter the only symptom is content quietly +// not federating. +// +// A COUNTER RATHER THAN A GAUGE, and it is a RATE that matters, not a total: +// the ordinary case (a native user voting in a native community) increments it +// constantly, so the number is meaningless in isolation and informative when it +// moves against its own baseline. +// +// Declared at package scope rather than inside PublishMetrics because the +// increment happens in the gate, which knows nothing about the connector the +// other gauges read; expvar.NewInt registers it exactly once per process. +var unclaimedSkips = expvar.NewInt(MetricUnclaimedSkips) + // deadLetterDepthUnavailable is what the backlog gauge reports when storage // cannot be read. A negative value is impossible for a count, so it is // unmistakable — where a 0 would claim the backlog is empty at exactly the diff --git a/internal/consume/profile.go b/internal/consume/profile.go index 979446e..0c5a13e 100644 --- a/internal/consume/profile.go +++ b/internal/consume/profile.go @@ -5,6 +5,7 @@ import ( "database/sql" "fmt" "log/slog" + "strconv" "tidepool/internal/errors" "tidepool/internal/store" @@ -24,6 +25,18 @@ import ( // A DID with no actor row is skipped rather than minted: a profile edit is not // a federating interaction, and minting here would give an AP identity to // every Coves user who ever set a display name. +// +// That skip CLAIMS its gate row, unlike the transient skips errSkipUnclaimed +// covers, and the trade is deliberate. The actor could exist tomorrow, so the +// skip is technically recoverable — but what is lost is one cached display +// name, the actor document already falls back to the local part, and the next +// profile write refreshes it. Against that: a profile edit by a user who never +// federates anything is one of the most common events on this stream, and +// releasing the claim would put every one of them in the unclaimed-skip counter +// and an INFO log, drowning the signal that counter exists to carry. +// +// The tx is unused for the same reason handleProfile writes no outbound state: +// this handler only refreshes a cache, and UpdateProfile is last-write-wins. func (d *Dispatcher) handleProfile(ctx context.Context, _ *sql.Tx, did string, commit *CommitEvent) error { if _, err := d.apActors.GetByDID(ctx, did); err != nil { if errors.IsNotFound(err) { @@ -69,14 +82,46 @@ func (d *Dispatcher) handleProfile(ctx context.Context, _ *sql.Tx, did string, c // The local part is untouched. It was frozen at actor creation, and // re-deriving it would strand every federated mention of the old name — a // rename may refresh the profile CACHE and nothing else. +// +// NO ORDERING GUARD, AND THAT IS THE RULING RATHER THAN AN OVERSIGHT. +// IdentityEvent.Seq is parsed off the wire and deliberately not consulted here, +// where the sibling #account tier claims its seq (applyAccountClaimed) before +// touching anything. The difference is what each handler WRITES. An #account +// frame's payload IS the state applied — active/status go straight into the +// actor row — so a stale replay writes stale facts and a user who came back +// silently stays paused. This handler writes nothing the frame carries: the +// handle is re-resolved from PLC and well-known every time, so a frame from an +// hour ago and a frame from a second ago apply the IDENTICAL current answer. +// Replaying one cannot regress the cache. +// +// The residual, named so a later change does not rediscover it as a surprise: +// two handlers racing for one DID could resolve in one order and write in the +// other, leaving the older handle cached. It is bounded to a single write — +// the fallback below fires only while DisplayName is empty, and the first +// write fills it — so the window closes after the first rename and no seq +// claim is worth holding an advisory lock across two network round-trips for. +// A seq claim becomes REQUIRED the moment this handler starts persisting +// anything the frame itself asserts. func (d *Dispatcher) handleIdentity(ctx context.Context, event *JetstreamEvent) error { if event.Identity == nil { return fmt.Errorf("%w: identity event for %s carries no identity", ErrPermanentEvent, event.DID) } - did := event.Identity.DID - if did == "" { - did = event.DID + + // The envelope DID is authoritative, exactly as it is for #account. A nested + // payload naming a DIFFERENT DID is malformed or hostile — acting on the + // inner one lets a frame about DID A resolve, verify and cache DID B — and + // no retry makes the two agree, so it is rejected as permanent BEFORE the + // resolver is reached. That last part is the security half: without it a + // crafted #identity frame aims this bridge's PLC and well-known round trips + // at any DID the attacker cares to name. + if event.Identity.DID != "" && event.Identity.DID != event.DID { + return fmt.Errorf("%w: identity payload DID %s disagrees with the envelope DID %s", + ErrPermanentEvent, strconv.Quote(event.Identity.DID), strconv.Quote(event.DID)) } + // An ABSENT inner DID is not a disagreement: the envelope answers for the + // frame. Jetstream always fills it, but the DLQ stores raw frames and a + // redriven or hand-repaired one can be thinner than the wire shape. + did := event.DID // The actor check comes FIRST: every Coves user who renames emits one of // these, and resolving would spend two network round-trips updating a diff --git a/internal/consume/profile_test.go b/internal/consume/profile_test.go index 2f19848..36053f1 100644 --- a/internal/consume/profile_test.go +++ b/internal/consume/profile_test.go @@ -188,6 +188,66 @@ func TestIdentityHandler_ActorlessDIDResolvesNothing(t *testing.T) { assert.Zero(t, countRows(t, database, "ap_actors"), "and no eager mint") } +// identityFrameSplitDID is an #identity frame whose NESTED payload names a +// different DID than the envelope — the exact disagreement handleAccount +// rejects as permanent. +func identityFrameSplitDID(envelopeDID, innerDID, handle string) []byte { + return []byte(fmt.Sprintf( + `{"did":%q,"time_us":6600,"kind":"identity",`+ + `"identity":{"did":%q,"handle":%q,"seq":10,"time":"2026-08-12T10:00:00.000Z"}}`, + envelopeDID, innerDID, handle)) +} + +// identityFrameNoInnerDID omits identity.did entirely. Jetstream always fills +// it, but the DLQ stores raw frames and a redriven or hand-repaired one can be +// thinner than the wire shape. +func identityFrameNoInnerDID(envelopeDID, handle string) []byte { + return []byte(fmt.Sprintf( + `{"did":%q,"time_us":6600,"kind":"identity",`+ + `"identity":{"handle":%q,"seq":10,"time":"2026-08-12T10:00:00.000Z"}}`, + envelopeDID, handle)) +} + +// Chunk 3 finding 5: handleIdentity trusted the NESTED payload's DID over the +// envelope's and never compared them, while its sibling handleAccount rejects +// exactly that disagreement. A crafted frame could therefore name any DID it +// liked and spend this bridge's PLC and well-known round-trips on it. +func TestIdentityHandler_InnerDIDDisagreeingWithTheEnvelopeIsPermanent(t *testing.T) { + database := dispatchTestDB(t) + const victimDID = "did:plc:44ybard66vv44zksje25o7dz" + seedAPActor(t, database, victimDID, "victim") + fixture := newDispatchFixture(t, database) + + err := fixture.handle(t, identityFrameSplitDID(dispatchNativeDID, victimDID, "attacker.example")) + + require.Error(t, err) + assert.ErrorIs(t, err, ErrPermanentEvent, + "a frame ABOUT DID A that carries DID B is malformed or hostile, and no retry "+ + "makes the two agree — the same ruling handleAccount already makes, for the "+ + "same reason: acting on the inner one lets a frame about A mutate B") + + assert.Empty(t, fixture.resolver.Calls(), + "and it is refused BEFORE the resolver runs: a crafted frame must not be able "+ + "to point this bridge's PLC and well-known lookups at any DID it names") + + displayName, _, _ := profileCache(t, database, victimDID) + assert.Empty(t, displayName, "no state is touched for the DID the payload named") +} + +func TestIdentityHandler_EmptyInnerDIDFallsBackToTheEnvelope(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + fixture := newDispatchFixture(t, database) + fixture.resolver.handle = "alice2.coves.social" + + require.NoError(t, fixture.handle(t, + identityFrameNoInnerDID(dispatchNativeDID, "alice2.coves.social")), + "an ABSENT inner DID is not a disagreement: the envelope answers for the frame") + + displayName, _, _ := profileCache(t, database, dispatchNativeDID) + assert.Equal(t, "alice2.coves.social", displayName) +} + func TestIdentityHandler_ResolverFailureLeavesTheCacheAlone(t *testing.T) { database := dispatchTestDB(t) seedAPActor(t, database, dispatchNativeDID, "alice") diff --git a/internal/consume/rev_gate.go b/internal/consume/rev_gate.go index 4f86311..b4676dd 100644 --- a/internal/consume/rev_gate.go +++ b/internal/consume/rev_gate.go @@ -118,6 +118,71 @@ func logSkippedStaleRev(consumer, operation, uri, rev string) { slog.String("rev", rev)) } +// errSkipUnclaimed is the sentinel a commit handler returns to skip an event +// WITHOUT claiming its gate row. applyGatedTx reports the event handled (the +// cursor advances) but rolls the claim back, so a redelivery of the same frame +// re-enters the handler instead of losing at `rev < EXCLUDED.rev`. +// +// WHICH SKIPS BELONG HERE. Only the ones whose answer can CHANGE while the +// event stays exactly as it is: +// +// - a vote or a comment on something this bridge has not materialized yet +// (the parent lands milliseconds later, and the vote arrived first); +// - a post for a community nobody has bridged yet (an operator bridges it); +// - a vote direction this build does not understand — `direction` is an OPEN +// enum, so the record is forward-compatible and a LATER BUILD is its +// recovery path. +// +// Everything else keeps claiming. An opted-out author's write, a delete with no +// outbound state, a malformed record: those are final for THIS frame, and a +// released claim would only re-run a decision that is already made — or, for a +// delete, throw away the tombstone that rejects the stale create. +// +// The first bullet is deliberately IMPRECISE, and the imprecision is the safe +// direction. resolveSubject answers nil for "never materialized" and for +// "materialized and since removed" alike, so a vote on a deleted post takes the +// unclaimed branch too. That costs a repeated skip if the frame is ever +// redelivered and nothing else — where guessing the other way would be the +// silent loss this sentinel exists to end. +// +// WHAT ACTUALLY REPLAYS THE EVENT, stated plainly because the sentinel is not a +// retry and must not be read as one. Jetstream delivers live traffic once, and +// a skip never dead-letters, so the DLQ redriver is NOT a path here. The real +// ones, in the order they matter: +// +// 1. the reconnect cursor REWIND. Every re-dial resumes a few seconds behind +// the last processed event (Connector.cursorRewind, 5s), which is exactly +// the window the parent-materializes-late race lives in. +// 2. a cursor behind Jetstream's retention, which replays the whole store. +// 3. a CursorSchemaVersion bump with the mandated truncate of +// jetstream_record_revs (GateResetRequiredOnSchemaBump) — the from-scratch +// replay a build that understands the new lexicon relies on. +// +// None of those is guaranteed to happen soon, and this sentinel does not make +// one happen. What it does is remove the POISON: with a gate row claimed, every +// one of those three replays is a silent no-op and the content is lost for +// good. Without it, each of them recovers the event. +var errSkipUnclaimed = stderrors.New("skipped without claiming the rev gate") + +// logSkipUnclaimed is the single, grep-able line for an unclaimed skip, and it +// is INFO rather than debug on purpose: unlike a stale-rev skip, this one means +// something did NOT happen that may still be owed. It is the only trace the +// event leaves anywhere — no DLQ row, no failed event, no state — so a +// deployment running at default level would otherwise have no way to see a +// backlog of votes waiting on a materializer that stalled. +// +// The reason travels in the error, not in an attribute, so each handler states +// its own case once and this stays the only place that formats it. +func logSkipUnclaimed(consumer, operation, uri, rev string, reason error) { + unclaimedSkips.Add(1) + slog.Info("rev-gate claim released: event skipped and left replayable", + slog.String("consumer", consumer), + slog.String("operation", operation), + slog.String("uri", uri), + slog.String("rev", rev), + slog.String("reason", reason.Error())) +} + // RevGate carries the gate's own DB handle. applyGated/applyGatedTx use it to // open the claim transaction held across apply, which is how EVERY commit // handler is currently gated. @@ -180,7 +245,8 @@ func (g *RevGate) Advance(ctx context.Context, uri, rev string) error { // gate row lock a pure per-record mutex around apply. applyGatedTx now hands // that transaction to handlers so their durable state commits with the gate // advance, and they use it: the comment path writes outbound_objects on it, and -// the federation path writes outbound_deliveries and ap_actors. +// the federation path writes federation_prefs, outbound_deliveries and +// ap_actors. // // The rule that replaces the old proof: A HANDLER WRITING ON THIS TRANSACTION // MUST NOT THEN CALL SOMETHING THAT OPENS A SECOND TRANSACTION TOUCHING THE @@ -188,14 +254,18 @@ func (g *RevGate) Advance(ctx context.Context, uri, rev string) error { // handler is synchronously waiting, so the two deadlock until a timeout. The // destructive opt-out tier is exactly that shape — the seam takes no // transaction and opens its own — which is why handleFederation defers its -// ap_actors write until after that call. +// ap_actors write until after that call, and why the same path is the ONE that +// commits its federation_prefs row ahead of this transaction instead of on it +// (applyOptOut): the purge marks that very row from its own transaction. // // An apply error — or a panic, which the deferred rollback covers equally — // releases the claim WITHOUT advancing, so the connector's retry/redrive // replays the event instead of losing it behind its own gate entry. // // A gate SKIP returns nil, not an error: the event is fully accounted for and -// the cursor must advance past it. +// the cursor must advance past it. So does an errSkipUnclaimed skip — with the +// difference that the claim is rolled back rather than committed, which is what +// keeps the event replayable. func applyGated(ctx context.Context, gate *RevGate, consumer, did string, commit *CommitEvent, apply func() error) error { return applyGatedTx(ctx, gate, consumer, did, commit, func(*sql.Tx) error { return apply() }) } @@ -218,10 +288,20 @@ func applyGated(ctx context.Context, gate *RevGate, consumer, did string, commit // get a validation error from UpsertTx, which is the correct refusal for a // path that has opted out of gating. func applyGatedTx(ctx context.Context, gate *RevGate, consumer, did string, commit *CommitEvent, apply func(tx *sql.Tx) error) error { + uri := commitRecordURI(did, commit) if gate == nil || commit.Rev == "" { - return apply(nil) + // The bypass still has to understand the sentinel: a handler cannot know + // whether it is running under a claim, and returning the skip as an error + // would dead-letter an event that was merely not-yet-applicable. + if err := apply(nil); err != nil { + if stderrors.Is(err, errSkipUnclaimed) { + logSkipUnclaimed(consumer, commit.Operation, uri, commit.Rev, err) + return nil + } + return err + } + return nil } - uri := commitRecordURI(did, commit) tx, err := gate.db.BeginTx(ctx, nil) if err != nil { return fmt.Errorf("begin rev-gate transaction for %s: %w", uri, err) @@ -242,6 +322,15 @@ func applyGatedTx(ctx context.Context, gate *RevGate, consumer, did string, comm return nil } if err := apply(tx); err != nil { + if stderrors.Is(err, errSkipUnclaimed) { + // Deliberately un-advanced: the deferred rollback releases the claim + // with no gate row written, so the same frame re-enters this handler + // on any of the replay paths errSkipUnclaimed documents. The event is + // still reported HANDLED — nothing is owed to the retry budget and + // the cursor must move past it. + logSkipUnclaimed(consumer, commit.Operation, uri, commit.Rev, err) + return nil + } return err // the deferred rollback releases the claim un-advanced } if err := tx.Commit(); err != nil { diff --git a/internal/consume/unclaimed_skip_test.go b/internal/consume/unclaimed_skip_test.go new file mode 100644 index 0000000..e7f5451 --- /dev/null +++ b/internal/consume/unclaimed_skip_test.go @@ -0,0 +1,222 @@ +package consume + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// Chunk 3 finding 6 (CRITICAL, silent failure): a handler skip INSIDE the rev +// gate commits the claim, and the retry that would fix it then loses at +// `rev < EXCLUDED.rev` and is logged at Debug as a stale replay. +// +// The user-visible loss: a vote or a comment that arrives a few hundred +// milliseconds before the thing it hangs on has materialized is skipped, the +// gate row is written anyway, and the byte-identical redelivery — the cursor +// rewind that exists precisely to recover this — is swallowed. No DLQ row, no +// metric, no log above Debug. +// +// These tests split the skips by whether the answer can CHANGE: +// +// - TRANSIENT skips must leave NO gate row (errSkipUnclaimed), so the same +// frame re-enters the handler on redelivery and applies. +// - PERMANENT skips must still claim, or a redelivery would re-execute a +// decision that is already final. + +// --------------------------------------------------------------------------- +// Transient skips: the gate row must NOT be written +// --------------------------------------------------------------------------- + +func TestUnclaimedSkip_VoteBeforeItsSubjectMaterializesAppliesOnRedelivery(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + // Deliberately NO seedThreadRoot: the subject this vote names has not been + // materialized yet, which is the ordinary few-hundred-millisecond race + // between a post being bridged and somebody voting on it. + ctx := context.Background() + + const rkey = "3lzskipvote01" + voteATURI := voteATURIFor(dispatchNativeDID, rkey) + frame := voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "up") + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, frame), + "an unresolved subject is a skip, not a failure: native users vote in native "+ + "communities constantly and dead-lettering that would bury the queue") + + assert.Empty(t, storedRev(t, database, voteATURI), + "NO gate row may be claimed for an event no handler applied. The claim is what "+ + "turns the redelivery into a silent stale-replay skip, which is how the "+ + "vote is lost") + assert.Zero(t, countRows(t, database, "outbound_votes")) + assert.Empty(t, fixture.enqueuer.Calls()) + + // The parent materializes, and the cursor rewind redelivers the SAME frame. + seedThreadRoot(t, database) + replay := newDispatchFixture(t, database) + require.NoError(t, replay.handle(t, frame)) + + stored, err := store.NewOutboundVotes(database).GetByATURI(ctx, voteATURI) + require.NoError(t, err, + "the byte-identical redelivery now applies: nothing was claimed the first time, "+ + "so the gate has nothing to reject it with") + require.NotNil(t, stored) + assert.Equal(t, "up", stored.Direction) + assert.Len(t, replay.enqueuer.Calls(), 1, "and the Like finally federates") + assert.Equal(t, dispatchRev, storedRev(t, database, voteATURI), + "and NOW the gate is claimed, so a further replay is the ordinary no-op") +} + +func TestUnclaimedSkip_CommentBeforeItsThreadMaterializesAppliesOnRedelivery(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + ctx := context.Background() + + const rkey = "3lzskipcmnt01" + atURI := commentATURIFor(dispatchNativeDID, rkey) + frame := commentFrameFull(dispatchNativeDID, dispatchRev, rkey, "create", "hi", acceptRootATURI) + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, frame)) + + assert.Empty(t, storedRev(t, database, atURI), + "a comment whose thread does not resolve YET must not claim its gate row") + assert.Zero(t, countRows(t, database, "outbound_objects")) + + seedThreadRoot(t, database) + replay := newDispatchFixture(t, database) + require.NoError(t, replay.handle(t, frame)) + + stored, err := store.NewOutboundObjects(database).GetByATURI(ctx, atURI) + require.NoError(t, err, "the redelivered comment federates once its thread exists") + require.NotNil(t, stored) + assert.Len(t, replay.enqueuer.Calls(), 1) +} + +func TestUnclaimedSkip_PostV2ForANotYetBridgedCommunityAppliesOnRedelivery(t *testing.T) { + database := dispatchTestDB(t) + // No community row: this community is not bridged at the moment the post + // arrives. Bridging one is an ordinary operator action that happens later. + + const rkey = "3lzskippost01" + postATURI := "at://" + dispatchNativeDID + "/" + CollectionPostV2 + "/" + rkey + frame := postV2Frame(dispatchNativeDID, dispatchRev, rkey, acceptCommunityDID) + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, frame)) + assert.Zero(t, fixture.engine.Calls()) + assert.Empty(t, storedRev(t, database, postATURI), + "a post for a community this bridge does not federate YET must not claim its "+ + "gate row: the community being bridged tomorrow is the whole recovery path") + + seedBridgedCommunity(t, database) + replay := newDispatchFixture(t, database) + require.NoError(t, replay.handle(t, frame)) + assert.Equal(t, 1, replay.engine.Calls(), + "once the community is bridged, the redelivered post reaches the acceptance engine") +} + +func TestUnclaimedSkip_UnrecognisedVoteDirectionStaysReplayable(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + + const rkey = "3lzskipdir001" + voteATURI := voteATURIFor(dispatchNativeDID, rkey) + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "sideways"))) + + assert.Zero(t, countRows(t, database, "outbound_votes"), + "nothing is stored: guessing a direction would push a vote the user never cast") + assert.Empty(t, storedRev(t, database, voteATURI), + "direction is an OPEN enum, so the record is forward-compatible rather than "+ + "malformed and a LATER BUILD is its recovery path. Claiming the gate here "+ + "would make the from-scratch replay that build depends on a silent no-op") +} + +// --------------------------------------------------------------------------- +// Permanent skips: the claim still commits +// --------------------------------------------------------------------------- + +func TestClaimedSkip_AnOptedOutAuthorsVoteStillClaimsTheGate(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + ctx := context.Background() + + _, err := store.NewFederationPrefs(database).Upsert(ctx, store.FederationPref{ + DID: dispatchNativeDID, + Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + + const rkey = "3lzskipout001" + voteATURI := voteATURIFor(dispatchNativeDID, rkey) + frame := voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "up") + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, frame)) + assert.Equal(t, dispatchRev, storedRev(t, database, voteATURI), + "PERMANENT: the author said no. The decision cannot become untrue for THIS "+ + "record, so the claim commits and the replay is rejected at the gate") + + // Even after the preference is gone, the identical frame is a gate skip: a + // re-enable federates what the user writes NEXT, not what they wrote while + // opted out. + require.NoError(t, store.NewFederationPrefs(database).Delete(ctx, dispatchNativeDID)) + replay := newDispatchFixture(t, database) + require.NoError(t, replay.handle(t, frame)) + assert.Zero(t, countRows(t, database, "outbound_votes")) + assert.Empty(t, replay.enqueuer.Calls(), + "the replay loses at `rev < EXCLUDED.rev` — which is the gate working") +} + +func TestClaimedSkip_AVoteDeleteWithNoStateKeepsItsTombstone(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + + const rkey = "3lzskipdel001" + voteATURI := voteATURIFor(dispatchNativeDID, rkey) + + // The delete arrives first (a rewind can reorder nothing, but a delete for a + // vote cast before this bridge existed looks exactly like this). + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, + voteDeleteFrame(dispatchNativeDID, dispatchRevHigher, rkey))) + assert.Equal(t, dispatchRevHigher, storedRev(t, database, voteATURI), + "PERMANENT, and load-bearing: the gate row a delete leaves IS the tombstone "+ + "that rejects a stale create. Releasing it unclaimed would let a replay "+ + "past the delete resurrect the vote") + + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "up"))) + assert.Zero(t, countRows(t, database, "outbound_votes"), + "and the stale create loses at the tombstone rather than re-casting the vote") +} + +// --------------------------------------------------------------------------- +// The counter +// --------------------------------------------------------------------------- + +func TestUnclaimedSkip_IsCounted(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + + before := unclaimedSkips.Value() + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, "3lzskipcnt01", acceptRootATURI, "up"))) + + assert.Equal(t, before+1, unclaimedSkips.Value(), + "a transient skip is the one outcome with NO other trace — no DLQ row, no "+ + "failed event — so the counter is the only thing that makes a stuck "+ + "backlog of them visible") +} diff --git a/internal/consume/votes.go b/internal/consume/votes.go index 1ba78a4..296a5d1 100644 --- a/internal/consume/votes.go +++ b/internal/consume/votes.go @@ -5,6 +5,7 @@ import ( "database/sql" "fmt" "log/slog" + "strconv" "tidepool/internal/errors" "tidepool/internal/store" @@ -37,6 +38,9 @@ func (d *Dispatcher) handleVote(ctx context.Context, tx *sql.Tx, did string, com case operationDelete: return d.applyVoteDelete(ctx, tx, did, commit) default: + // Unreachable: validateCommitEnvelope rejects any other operation as + // permanent before the gate. Kept as a claimed skip because an operation + // this build cannot name is not something a replay improves. d.logger.Debug("unknown vote operation", slog.String("operation", commit.Operation), slog.String("did", did)) return nil @@ -56,6 +60,9 @@ func (d *Dispatcher) applyVoteWrite(ctx context.Context, tx *sql.Tx, did string, return err } if !federating { + // PERMANENT, and it CLAIMS the gate. The author asked us to stop, which + // is not a fact that changes underneath this frame: a later re-enable + // federates what they write next, not what they wrote while opted out. d.logger.Debug("skipping vote from an opted-out author", slog.String("did", did), slog.String("rkey", commit.RKey)) return nil @@ -66,13 +73,20 @@ func (d *Dispatcher) applyVoteWrite(ctx context.Context, tx *sql.Tx, did string, // Nothing is stored: guessing a direction would push a vote the user // never cast, and dead-lettering would turn a lexicon rollout into a // queue full of rows nobody can redrive. - d.logger.Debug("skipping vote with an unrecognised direction", - slog.String("did", did), slog.String("direction", direction)) - return nil + // + // UNCLAIMED, because `direction` is an OPEN enum: this record is + // forward-compatible rather than malformed, and the recovery path is a + // LATER BUILD replaying it (see errSkipUnclaimed). A claimed gate row + // would make that replay a silent no-op, which is the one outcome a + // forward-compatible field must not produce. + return fmt.Errorf("%w: vote %s has direction %s, which this build does not understand", + errSkipUnclaimed, commitRecordURI(did, commit), strconv.Quote(direction)) } subjectATURI := refURI(commit.Record, "subject") if subjectATURI == "" { + // PERMANENT: the lexicon requires subject, and no later state makes an + // absent field appear. The claim commits. d.logger.Debug("skipping vote with no subject", slog.String("did", did)) return nil } @@ -83,9 +97,13 @@ func (d *Dispatcher) applyVoteWrite(ctx context.Context, tx *sql.Tx, did string, if subject == nil { // Native users vote in native communities constantly; dead-lettering // that would bury the queue. - d.logger.Debug("skipping vote on a subject this bridge does not federate", - slog.String("did", did), slog.String("subject", subjectATURI)) - return nil + // + // UNCLAIMED. This is the finding's headline case: a vote cast a few + // hundred milliseconds before its subject was materialized resolves to + // nothing, and a claimed gate row would swallow the reconnect rewind + // that exists to recover it — the vote would silently never federate. + return fmt.Errorf("%w: vote %s names subject %s, which this bridge does not federate (yet)", + errSkipUnclaimed, commitRecordURI(did, commit), strconv.Quote(subjectATURI)) } if err := d.refuseBannedAuthor(ctx, did, subject.CommunityDID, "vote "+commitRecordURI(did, commit)); err != nil { @@ -164,6 +182,11 @@ func (d *Dispatcher) applyVoteDelete(ctx context.Context, tx *sql.Tx, did string stored, err := d.votes.GetByATURI(ctx, voteATURI) if errors.IsNotFound(err) { // No Undo may be sent for a Like no peer ever received. + // + // PERMANENT, and the claim is LOAD-BEARING rather than merely harmless: + // the gate row a delete leaves IS the tombstone that rejects a stale + // create for the same URI. Releasing it would let a replay past the + // delete re-cast a vote the user withdrew. d.logger.Debug("skipping delete for a vote with no outbound state", slog.String("did", did), slog.String("vote", voteATURI)) return nil diff --git a/internal/store/federation_prefs.go b/internal/store/federation_prefs.go index 7722139..600318a 100644 --- a/internal/store/federation_prefs.go +++ b/internal/store/federation_prefs.go @@ -21,6 +21,17 @@ func NewFederationPrefs(db *sql.DB) FederationPrefs { const federationPrefColumns = `did, enabled, delete_remote, source, purged_at, updated_at` func (r *postgresFederationPrefs) Upsert(ctx context.Context, pref FederationPref) (*FederationPref, error) { + return r.upsert(ctx, r.db, pref) +} + +func (r *postgresFederationPrefs) UpsertTx(ctx context.Context, tx *sql.Tx, pref FederationPref) (*FederationPref, error) { + if tx == nil { + return nil, errors.NewValidationError("tx", "must not be nil") + } + return r.upsert(ctx, tx, pref) +} + +func (r *postgresFederationPrefs) upsert(ctx context.Context, ex execer, pref FederationPref) (*FederationPref, error) { if !pref.Source.Valid() { // Including the zero value: an unstated source hides whether the // preference came off a record we saw or a probe we made, and those @@ -48,7 +59,7 @@ func (r *postgresFederationPrefs) Upsert(ctx context.Context, pref FederationPre updated_at = now() RETURNING ` + federationPrefColumns - row := r.db.QueryRowContext(ctx, query, + row := ex.QueryRowContext(ctx, query, pref.DID, pref.Enabled, pref.DeleteRemote, string(pref.Source)) stored, err := scanFederationPref(row) if err != nil { @@ -84,11 +95,28 @@ func (r *postgresFederationPrefs) Delete(ctx context.Context, did string) error // which is the most innocuous-looking way to undo an irreversible decision. // The caller is told nothing changed by reading the row back, which // restoreDefaultFederation does before it says anything to an operator. - if _, err := r.db.ExecContext(ctx, - `DELETE FROM federation_prefs WHERE did = $1 AND purged_at IS NULL`, did); err != nil { - return fmt.Errorf("delete federation_pref %q: %w", did, err) + _, err := deleteFederationPref(ctx, r.db, did) + return err +} + +func (r *postgresFederationPrefs) DeleteTx(ctx context.Context, tx *sql.Tx, did string) (bool, error) { + if tx == nil { + return false, errors.NewValidationError("tx", "must not be nil") } - return nil + return deleteFederationPref(ctx, tx, did) +} + +func deleteFederationPref(ctx context.Context, ex execer, did string) (bool, error) { + result, err := ex.ExecContext(ctx, + `DELETE FROM federation_prefs WHERE did = $1 AND purged_at IS NULL`, did) + if err != nil { + return false, fmt.Errorf("delete federation_pref %q: %w", did, err) + } + affected, err := result.RowsAffected() + if err != nil { + return false, fmt.Errorf("delete federation_pref %q: rows affected: %w", did, err) + } + return affected > 0, nil } func (r *postgresFederationPrefs) MarkPurgedTx(ctx context.Context, tx *sql.Tx, did string) error { diff --git a/internal/store/federation_prefs_test.go b/internal/store/federation_prefs_test.go index cef149b..bbb7a97 100644 --- a/internal/store/federation_prefs_test.go +++ b/internal/store/federation_prefs_test.go @@ -251,3 +251,108 @@ func TestFederationPrefs_UpsertNeverTouchesPurgedAt(t *testing.T) { require.NotNil(t, got.PurgedAt, "and the read-back agrees with the returned row") assert.True(t, got.PurgedAt.Equal(*purged.PurgedAt)) } + +// --------------------------------------------------------------------------- +// The transactional flavors (chunk 3 finding 4) +// --------------------------------------------------------------------------- +// +// The consumer's opt-out door writes its preference on the rev-gate's +// transaction so the preference, the delivery cancellation and the gate advance +// commit as one unit. These pin the property that makes that possible — that +// these two really do ride the caller's transaction — because a method that +// quietly ran on r.db instead would leave every consumer-side atomicity test +// passing for the wrong reason. + +func TestFederationPrefs_UpsertTxRidesTheCallersTransaction(t *testing.T) { + database := federationPrefsTestDB(t) + repo := NewFederationPrefs(database) + ctx := context.Background() + + tx, err := database.BeginTx(ctx, nil) + require.NoError(t, err) + stored, err := repo.UpsertTx(ctx, tx, fpOptOut(testDID, FederationPrefSourceRecord)) + require.NoError(t, err) + require.NotNil(t, stored) + require.NoError(t, tx.Rollback()) + + _, err = repo.Get(ctx, testDID) + require.Error(t, err) + assert.True(t, errors.IsNotFound(err), + "a rolled-back transaction leaves NO preference: the write is the caller's to "+ + "commit, which is what lets the opt-out roll back with its gate advance") + + // Positive control: the same call on a committed transaction does write. + tx, err = database.BeginTx(ctx, nil) + require.NoError(t, err) + _, err = repo.UpsertTx(ctx, tx, fpOptOut(testDID, FederationPrefSourceRecord)) + require.NoError(t, err) + require.NoError(t, tx.Commit()) + + got, err := repo.Get(ctx, testDID) + require.NoError(t, err) + assert.False(t, got.Enabled) +} + +func TestFederationPrefs_DeleteTxRidesTheCallersTransactionAndReportsTheOutcome(t *testing.T) { + database := federationPrefsTestDB(t) + repo := NewFederationPrefs(database) + ctx := context.Background() + + _, err := repo.Upsert(ctx, fpOptOut(testDID, FederationPrefSourceRecord)) + require.NoError(t, err) + + tx, err := database.BeginTx(ctx, nil) + require.NoError(t, err) + deleted, err := repo.DeleteTx(ctx, tx, testDID) + require.NoError(t, err) + assert.True(t, deleted, "a standing preference is cleared, and the caller is told so") + require.NoError(t, tx.Rollback()) + + _, err = repo.Get(ctx, testDID) + require.NoError(t, err, + "and the rollback restores it: absence means default-on, so a delete that "+ + "outlived a failed event would silently re-enable a user who asked us to stop") + + // A COMMITTED PURGE is refused, and the false is how the caller learns that + // this is not an ordinary re-enable. + purged := mustPurgedPref(t, repo, testSecondDID, FederationPrefSourceRecord) + require.NotNil(t, purged.PurgedAt) + + tx, err = database.BeginTx(ctx, nil) + require.NoError(t, err) + deleted, err = repo.DeleteTx(ctx, tx, testSecondDID) + require.NoError(t, err) + assert.False(t, deleted, + "a withdrawn identity has nothing to come back to, and the caller must be able "+ + "to tell that apart from 'there was nothing to clear'") + require.NoError(t, tx.Commit()) + + survivor, err := repo.Get(ctx, testSecondDID) + require.NoError(t, err, "the purged row survives its own delete") + require.NotNil(t, survivor.PurgedAt) + + // And a DID that never opted out is also false — same signal, different + // meaning, which is why restoreDefaultFederation reads the row back. + tx, err = database.BeginTx(ctx, nil) + require.NoError(t, err) + deleted, err = repo.DeleteTx(ctx, tx, fpThirdDID) + require.NoError(t, err) + assert.False(t, deleted) + require.NoError(t, tx.Commit()) +} + +func TestFederationPrefs_TxFlavorsRefuseANilTransaction(t *testing.T) { + database := federationPrefsTestDB(t) + repo := NewFederationPrefs(database) + ctx := context.Background() + + _, err := repo.UpsertTx(ctx, nil, fpOptOut(testDID, FederationPrefSourceRecord)) + require.Error(t, err) + assert.True(t, errors.IsValidation(err), + "a nil tx is the correct refusal for a path that has opted out of the gate, not "+ + "a silent fall-back to autocommit") + + _, err = repo.DeleteTx(ctx, nil, testDID) + require.Error(t, err) + assert.True(t, errors.IsValidation(err)) +} diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index a09cc7d..693d0e6 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -613,6 +613,16 @@ type FederationPrefs interface { // the zero value is an error satisfying errors.IsValidation. Upsert(ctx context.Context, pref FederationPref) (*FederationPref, error) + // UpsertTx is Upsert on an existing transaction — the seam the consumer's + // opt-out handler uses so the preference, the delivery cancellation and the + // rev-gate advance commit as ONE unit. A nil tx is an error satisfying + // errors.IsValidation. + // + // NOT for a caller that then reaches a seam opening its own transaction + // against this row: it cannot release what it holds without committing, and + // the two block each other (see consume/rev_gate.go's DEADLOCK NOTE). + UpsertTx(ctx context.Context, tx *sql.Tx, pref FederationPref) (*FederationPref, error) + // Get returns the preference for a DID. A miss is an error satisfying // errors.IsNotFound and MEANS default-on, not "unknown". Get(ctx context.Context, did string) (*FederationPref, error) @@ -624,6 +634,19 @@ type FederationPrefs interface { // come back to. Callers that report an outcome read the row back. Delete(ctx context.Context, did string) error + // DeleteTx is Delete on an existing transaction — the re-enable half of the + // consumer's opt-out door, so the clearing rides the rev-gate advance and a + // failed event cannot leave federation silently restored for a user who + // asked us to stop. A nil tx is an error satisfying errors.IsValidation. + // + // It reports whether a row was actually removed, which is how a caller tells + // the ordinary re-enable from the one case this refuses: a COMMITTED PURGE + // is not deletable. A false with no error means either "there was nothing to + // clear" or "the identity was withdrawn and cannot come back", and since + // nothing was written, the caller may read the row back on any connection to + // tell which. + DeleteTx(ctx context.Context, tx *sql.Tx, did string) (deleted bool, err error) + // MarkPurged records that the destructive tier actually asked peers to // delete this user's content — the fact that makes the preference terminal. // The FIRST commit wins; a retry never moves the date. A missing preference