diff --git a/internal/consume/account_confirm_test.go b/internal/consume/account_confirm_test.go new file mode 100644 index 0000000..ba0b557 --- /dev/null +++ b/internal/consume/account_confirm_test.go @@ -0,0 +1,288 @@ +package consume + +import ( + "context" + "database/sql" + stderrors "errors" + "io" + "log/slog" + "sync" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/optout" + "tidepool/internal/store" +) + +// TASK 17d, CYCLE 4 — THE FOUR OUTCOMES OF THE CONFIRM (decision 19). +// +// A #account frame is a CLAIM about a moment that may have passed: the reconnect +// rewind replays it, a user can reactivate, and what was true an hour ago need +// not be true now. The action it proposes — asking every peer to purge a user's +// content — is one no peer undoes. So the terminal tier confirms against the +// identity's own sources first, and the confirm has THREE answers, not two: +// +// deleted → act; live → do nothing, and be DONE with the event; +// error → we do not know, which is not "not deleted". +// +// Collapsing the third into either of the others is the bug the seam exists to +// prevent, and it is 17c-3's P1-c again (absent and unparseable both becoming +// nil, so an unreadable expiry meant a permanent ban): read as live it silently +// drops a real deletion and leaves the user federated forever; read as deleted it +// erases someone who never left. +// +// THE SEQ IS THE OTHER HALF, and it is what separates "live" from "unknown" +// operationally. #account is ungated by the rev gate — its own per-DID seq is +// the ordering guard — so an outcome that does not advance it makes the event +// redrivable, and one that does closes it. A stale deletion for a live account +// must NOT wedge every later status change for that DID behind itself; an +// unconfirmable one must not be silently consumed. +// +// The confirmer here is a stub on purpose: the real PLC→#atproto_pds→ +// getRepoStatus resolver is covered where it lives. What is under test is what +// the tier DOES with each of the three answers. + +// stubConfirmer is the AccountConfirmer seam, driven to each of its three +// outcomes. It counts calls so "nothing destructive happened" can be told apart +// from "the tier was never reached". +type stubConfirmer struct { + mu sync.Mutex + deleted bool + err error + calls []string +} + +func (c *stubConfirmer) AccountStatus(_ context.Context, did string) (bool, error) { + c.mu.Lock() + defer c.mu.Unlock() + c.calls = append(c.calls, did) + return c.deleted, c.err +} + +func (c *stubConfirmer) Calls() []string { + c.mu.Lock() + defer c.mu.Unlock() + return append([]string(nil), c.calls...) +} + +func (c *stubConfirmer) set(deleted bool, err error) { + c.mu.Lock() + defer c.mu.Unlock() + c.deleted, c.err = deleted, err +} + +// terminalWorld is a dispatcher whose terminal tier is the REAL +// optout.Terminator — the thing that owns the confirm-then-act sequence — +// wired to a stub confirmer and a recording destructive seam. +type terminalWorld struct { + *dispatchFixture + confirmer *stubConfirmer + purger *recordingDeleter +} + +// newTerminalWorld builds that world. deleterWired=false is the deployment +// where the destructive tier has not landed: the Terminator's Deleter is nil, +// which is a different thing from a deleter that is never called. +func newTerminalWorld(t *testing.T, database *sql.DB, deleterWired bool) *terminalWorld { + t.Helper() + confirmer := &stubConfirmer{} + purger := &recordingDeleter{} + opts := optout.Options{ + Confirmer: confirmer, + Prefs: store.NewFederationPrefs(database), + // Discarded because NewTerminator warns at CONSTRUCTION when the seam is + // absent: a log assertion here would pass without the handler ever + // announcing anything. + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + } + if deleterWired { + opts.Deleter = purger + } + terminator, err := optout.NewTerminator(opts) + require.NoError(t, err) + + fixture := newDispatchFixture(t, database, + func(o *Options) { o.Terminator = terminator }) + return &terminalWorld{dispatchFixture: fixture, confirmer: confirmer, purger: purger} +} + +// appliedSeq is the last #account seq the consumer recorded for a DID. It is +// the cursor for that DID's status history: an event that does not advance it +// is redelivered, and one that does is closed. +func appliedSeq(t *testing.T, database *sql.DB, did string) int64 { + t.Helper() + var seq int64 + require.NoError(t, database.QueryRowContext(context.Background(), + `SELECT last_account_seq FROM ap_actors WHERE did = $1`, did).Scan(&seq)) + return seq +} + +// storedPref reads the terminal preference, or nil when none was recorded. +func storedPref(t *testing.T, database *sql.DB, did string) *store.FederationPref { + t.Helper() + pref, err := store.NewFederationPrefs(database).Get(context.Background(), did) + if err != nil { + return nil + } + return pref +} + +// --------------------------------------------------------------------------- +// (1) CONFIRMED DELETED — the destructive path runs +// --------------------------------------------------------------------------- + +func TestConfirmedDeletion_RunsTheDestructivePathAndClosesTheEvent(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + world := newTerminalWorld(t, database, true) + world.confirmer.set(true, nil) + + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 11))) + + require.Equal(t, []string{dispatchNativeDID}, world.confirmer.Calls(), + "the event is never acted on directly: the DID's own sources are asked first, "+ + "because a deletion that was true when the frame was written may have been "+ + "reversed by the time we replay it") + assert.Equal(t, []string{dispatchNativeDID}, world.purger.DIDs(), + "and a CONFIRMED deletion runs the destructive tier — a user who is gone must not "+ + "keep a federated identity speaking for them on every instance that holds "+ + "their content") + + pref := storedPref(t, database, dispatchNativeDID) + require.NotNil(t, pref, + "the terminal state is RECORDED: peers that honour a Delete cannot restore what "+ + "they dropped, so the decision has to survive a crash between deciding and sending") + assert.False(t, pref.Enabled) + assert.True(t, pref.DeleteRemote, "recorded as the destructive tier, not as a soft opt-out") + assert.Equal(t, store.FederationPrefSourceAccount, pref.Source, + "and attributed to the ACCOUNT: the user wrote no record and there is none left to "+ + "fetch, so this column is the only place an operator can tell a deletion from "+ + "an opt-out") + + assert.Equal(t, int64(11), appliedSeq(t, database, dispatchNativeDID), + "the event is closed: a handled deletion must not be replayed into a second "+ + "purge on the next reconnect") +} + +// --------------------------------------------------------------------------- +// (2) CONFIRMED STILL LIVE — nothing destructive, and the seq STILL ADVANCES +// --------------------------------------------------------------------------- + +func TestUnconfirmedDeletionOfALiveAccountDoesNothingAndStillAdvances(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + world := newTerminalWorld(t, database, true) + // The frame says deleted; the identity says otherwise. This is the ordinary + // shape of a rewound cursor replaying a deletion a user has since reversed. + world.confirmer.set(false, nil) + + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 12)), + "a divergence between the firehose and the identity is not a failure: it is the "+ + "answer the confirm exists to get") + + require.NotEmpty(t, world.confirmer.Calls(), "precondition: the confirm really ran") + assert.Empty(t, world.purger.DIDs(), + "NOTHING destructive: an unconfirmed deletion that purged anyway would erase a "+ + "living user's content from every instance holding it, on the strength of an "+ + "event we just proved wrong") + assert.Nil(t, storedPref(t, database, dispatchNativeDID), + "and nothing is recorded either — a preference row here would stop federating for "+ + "a user who never asked, with no record left to contradict it") + + paused := deliveryPaused(t, database, dispatchNativeDID) + assert.False(t, paused, "the live account keeps delivering") + enabled, _ := actorEnabled(t, database, dispatchNativeDID) + assert.True(t, enabled, "under its original identity") + + assert.Equal(t, int64(12), appliedSeq(t, database, dispatchNativeDID), + "AND THE SEQ ADVANCES. Leaving it un-advanced would hold a stale deletion open "+ + "forever: the same frame redelivers on every reconnect, and — because the seq "+ + "gate rejects anything at or below the last APPLIED one — every later status "+ + "change for this DID queues behind an event that can never succeed") +} + +// --------------------------------------------------------------------------- +// (3) CONFIRM FAILED — retryable, nothing sent, seq NOT advanced +// --------------------------------------------------------------------------- + +// TestUnconfirmableDeletionIsRetryableAndAdvancesNothing is the case the seam +// exists for. "We could not confirm" must not collapse into either verdict: the +// event has to come back, which means it must NOT be closed. +func TestUnconfirmableDeletionIsRetryableAndAdvancesNothing(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + world := newTerminalWorld(t, database, true) + world.confirmer.set(false, stderrors.New("plc directory unreachable")) + + err := world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 13)) + require.Error(t, err, + "an unconfirmable deletion is an ERROR, which is what keeps the event redrivable: "+ + "returning nil would consume the user's deletion because a directory was down "+ + "for a minute") + assert.NotErrorIs(t, err, ErrPermanentEvent, + "and a TRANSIENT one: dead-lettering it exhausted turns a network blip into a "+ + "deletion that never happens") + + assert.Empty(t, world.purger.DIDs(), + "nothing is sent on an unknown: acting on a stale deletion is the unrecoverable "+ + "direction, and it is unrecoverable at every peer at once") + assert.Nil(t, storedPref(t, database, dispatchNativeDID), + "and nothing is recorded, so a later confirm decides on the evidence rather than "+ + "on a half-written state") + assert.Zero(t, appliedSeq(t, database, dispatchNativeDID), + "the seq must NOT advance: it is the only thing that brings this event back, and a "+ + "deletion consumed by an advance is gone — the user stays federated forever "+ + "with nothing left to notice it") + + // And the retry is a real one: the same frame, once the directory answers. + world.confirmer.set(true, nil) + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 13)), + "the redelivered frame is not rejected as stale — an un-advanced seq is what makes "+ + "the redrive possible") + assert.Equal(t, []string{dispatchNativeDID}, world.purger.DIDs(), + "and now it acts, on a confirmation instead of an assumption") + assert.Equal(t, int64(13), appliedSeq(t, database, dispatchNativeDID)) +} + +// --------------------------------------------------------------------------- +// (4) CONFIRMED DELETED, NO DESTRUCTIVE SEAM — recorded, nothing sent +// --------------------------------------------------------------------------- + +// TestConfirmedDeletionWithNoDestructiveSeamRecordsTheBacklogAndSendsNothing +// pins the deployment where the tier has not landed. It matches what +// deleteRemote=true already does with no seam wired — record the intent, leave +// the work in the backlog — rather than degrading into a pause, which looks +// handled in the database and is not. +func TestConfirmedDeletionWithNoDestructiveSeamRecordsTheBacklogAndSendsNothing(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + world := newTerminalWorld(t, database, false) // no Deleter on the Terminator + world.confirmer.set(true, nil) + + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 14)), + "a missing destructive seam must not panic and must not fail the event") + + pref := storedPref(t, database, dispatchNativeDID) + require.NotNil(t, pref, + "the preference is recorded even with nowhere to send it, so the deletion can be "+ + "acted on from the backlog when the tier lands instead of being lost") + assert.False(t, pref.Enabled) + assert.True(t, pref.DeleteRemote) + assert.Equal(t, store.FederationPrefSourceAccount, pref.Source) + + assert.Zero(t, countRows(t, database, "outbound_activities"), + "and NOTHING is sent: there is no seam to send it with, and a tier that half-acts "+ + "is worse than one that has not landed") + assert.Zero(t, countRows(t, database, "outbound_deliveries")) + + assert.False(t, deliveryPaused(t, database, dispatchNativeDID), + "it is NOT degraded into a pause: a pause says 'this user is coming back', which is "+ + "the one thing a confirmed deletion rules out, and it would satisfy an operator "+ + "looking for evidence the deletion was handled") + + assert.Equal(t, int64(14), appliedSeq(t, database, dispatchNativeDID), + "the event is closed on the strength of the RECORD: the backlog is durable, so "+ + "replaying the frame forever adds nothing") +} diff --git a/internal/consume/dispatch_test.go b/internal/consume/dispatch_test.go index 444ee1f..bf34dbd 100644 --- a/internal/consume/dispatch_test.go +++ b/internal/consume/dispatch_test.go @@ -37,7 +37,12 @@ func dispatchTestDB(t *testing.T) *sql.DB { // here the moment any test WRITES it, or a leftover row refuses the // next test's author and the failure names neither the ban nor the // test that left one. - "object_moderation", "community_bans") + "object_moderation", "community_bans", + // The outbound queue. Nothing in this package writes it through a + // recorder, which is exactly why it belongs here: 17d's terminal tier + // asserts that an unwired destructive seam sends NOTHING, and "no rows" + // is only an assertion about this event if the table started empty. + "outbound_deliveries", "outbound_activities") return database } diff --git a/internal/ingest/moderation_terminal_test.go b/internal/ingest/moderation_terminal_test.go index a4c1cb7..0ea5804 100644 --- a/internal/ingest/moderation_terminal_test.go +++ b/internal/ingest/moderation_terminal_test.go @@ -91,7 +91,12 @@ type moderationWorld struct { digestRKey string } -func newModerationWorld(t *testing.T, h *harness) moderationWorld { +// newModerationWorld builds that world. The variadic hooks mutate the +// consumer's Options just before the dispatcher is built, which is the only way +// to express a deployment where one 17d SEAM IS ABSENT — a nil seam and a seam +// that is never called are different states, and only the first one can be +// reached by leaving the field unset. +func newModerationWorld(t *testing.T, h *harness, mutate ...func(*consume.Options)) moderationWorld { t.Helper() ctx := context.Background() groupA := h.subscribeTechnology() @@ -149,7 +154,7 @@ func newModerationWorld(t *testing.T, h *harness) moderationWorld { UserOrigin: mtUserOrigin, }) require.NoError(t, err) - dispatcher, err := consume.NewDispatcher(consume.Options{ + consumerOpts := consume.Options{ DB: h.db, Actors: userOrigin, Enqueuer: enqueuer, @@ -161,7 +166,11 @@ func newModerationWorld(t *testing.T, h *harness) moderationWorld { // never reached. RemoteDeleter: outbound.NewPurger(h.db, mtUserOrigin, enqueuer), UserOrigin: mtUserOrigin, - }) + } + for _, apply := range mutate { + apply(&consumerOpts) + } + dispatcher, err := consume.NewDispatcher(consumerOpts) require.NoError(t, err) // --- The post is admitted through the real path: acceptance record, diff --git a/internal/ingest/optout_controls_test.go b/internal/ingest/optout_controls_test.go new file mode 100644 index 0000000..f9d308d --- /dev/null +++ b/internal/ingest/optout_controls_test.go @@ -0,0 +1,169 @@ +package ingest + +import ( + "context" + "database/sql" + "net/http" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/consume" + "tidepool/internal/store" +) + +// TASK 17d, CYCLE 5 — THE TWO CONTROLS. +// +// Both are about the tier NOT reached, which is the half of a two-tier design +// nothing else measures. A soft opt-out that quietly minted an identity, or a +// destructive one that half-acted with no seam to act through, would each pass +// every assertion the tiers make about themselves — the damage is entirely in +// what happened BESIDES the thing under test. + +const ( + // A DID that has never federated anything: no actor, and deliberately no + // entry in mtHandles either. The resolver refuses an unregistered DID, so an + // implementation that tried to mint here fails LOUDLY at the handle lookup + // rather than quietly creating the identity — the control's own tripwire. + ocStrangerDID = "did:plc:ocstranger000001" + ocStrangerRev = "3lzocrev000001" + + ocUnwiredRev = "3lzocrev000002" + ocUnwiredRKey = "3lzocpost00001" + ocUnwiredBRKey = "3lzocpost00002" +) + +// TestOptingOutFromADIDWithNoActorIsANoOpSuccess is the control the DEFAULT-ON +// design makes necessary. +// +// Under opt-out, the users who write this record are disproportionately the ones +// who have never federated anything — they read the disclosure and said no +// before their first bridged interaction. There is nothing to disable, and +// minting an actor in order to disable it would create the very identity the +// record asks the bridge not to create: a fediverse persona, discoverable and +// resolvable, brought into existence by the user's refusal. +// +// handleAccount already treats an actorless DID this way for #account. This +// pins it for the federation record, where the incentive to "make the mirror +// write succeed" is strongest. +func TestOptingOutFromADIDWithNoActorIsANoOpSuccess(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + actorsBefore := rowCount(t, h.db, "outbound_activities") + require.Zero(t, actorCount(t, h.db, ocStrangerDID), + "precondition: this DID has no AP identity — that is the whole case") + + // --- WHEN: they opt out, having never federated anything. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, ocStrangerDID, ocStrangerRev, "create", false, false, 1_775_000_040_000_001)), + "an opt-out for a DID with no actor is a no-op SUCCESS: the mirror write is allowed "+ + "to find nothing, and failing the event would dead-letter the most common opt-out "+ + "there is under a default-on design") + + // --- THEN: no identity was brought into existence to disable. + assert.Zero(t, actorCount(t, h.db, ocStrangerDID), + "NO actor may appear: a persona minted here is discoverable, resolvable and "+ + "webfinger-able — the bridge would have published a fediverse identity for a user "+ + "whose only instruction was not to") + + // --- AND: the preference is still recorded, because federation_prefs is the + // authority and answers for DIDs that have no actor to mirror onto. + pref, err := store.NewFederationPrefs(h.db).Get(ctx, ocStrangerDID) + require.NoError(t, err, + "the preference IS recorded: it is what gates the user's first federating "+ + "interaction, and an opt-out that stored nothing would be honoured only until "+ + "they replied to something") + assert.False(t, pref.Enabled) + assert.False(t, pref.DeleteRemote, "nothing destructive is ever inferred from a soft opt-out") + + assert.Equal(t, actorsBefore, rowCount(t, h.db, "outbound_activities"), + "and nothing was said about them on the wire: there is no identity to withdraw and "+ + "no peer that has heard of them") +} + +// TestDeleteRemoteWithNoDestructiveSeamRecordsWithoutActing is the control for +// the deployment where the destructive tier is not wired. +// +// The requirement is that it behaves like a BACKLOG, not like a degraded tier. +// The preference is durable, so the erasure can be carried out when the seam +// lands; and nothing is half-done in the meantime — because every partial step +// is a lie told to a different reader. A tombstoned actor with no Delete{Person} +// sent tells peers to stop resolving an identity whose content nobody was asked +// to purge; a vote flipped to 'undone' with no Undo on the wire hides that vote +// from the enumeration the real purge will one day run, so the erasure MISSES it. +func TestDeleteRemoteWithNoDestructiveSeamRecordsWithoutActing(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + // A deployment where the destructive tier has not landed. + world := newModerationWorld(t, h, func(o *consume.Options) { o.RemoteDeleter = nil }) + + // --- GIVEN: an author with real federated content and a vote a peer still + // holds. Without the live vote, "no Undo was enqueued" is true of a + // fixture that had nothing to undo, which is the same sentence about + // nothing. + admitPost(t, world, mtAuthorDID, ocUnwiredRKey, world.communityADID, "3lzocrev000010", 1_775_000_041_000_001) + admitPost(t, world, mtAuthorDID, ocUnwiredBRKey, world.communityBDID, "3lzocrev000011", 1_775_000_041_000_002) + requireEveryDelivery(t, h.db, mtAuthorDID, groupID, "pending", + "precondition: the author has queued work to withdraw") + + liveVote := seedDeliveredVote(t, h.db, mtAuthorDID, mtPostATURI, world.communityADID, "delivered") + require.Equal(t, "delivered", voteState(t, h.db, liveVote), + "precondition: a peer really is holding one of this actor's votes") + + activitiesBefore := rowCount(t, h.db, "outbound_activities") + + // --- WHEN: they ask for the destructive tier and there is nowhere to send it. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, mtAuthorDID, ocUnwiredRev, "create", false, true, 1_775_000_042_000_001)), + "a missing destructive seam must not panic and must not fail the event: dead-lettering "+ + "it would lose the user's request entirely") + + // --- THEN: the request is RECORDED, in the tier the user asked for. + pref, err := store.NewFederationPrefs(h.db).Get(ctx, mtAuthorDID) + require.NoError(t, err) + assert.False(t, pref.Enabled) + assert.True(t, pref.DeleteRemote, + "recorded as DESTRUCTIVE, not downgraded to the soft tier: the difference is the "+ + "user's own choice, and a preference that forgot it would leave their content "+ + "federated forever while the database says they were handled") + + // --- AND: the soft tier still ran, which is what makes the assertions below + // about restraint rather than about an event that did nothing. + actor, err := store.NewAPActors(h.db).GetByDID(ctx, mtAuthorDID) + require.NoError(t, err) + assert.False(t, actor.Enabled, "the actor is disabled: everything reachable was still done") + assertEveryDelivery(t, h.db, mtAuthorDID, groupID, "cancelled", + "and their queued work is cancelled — the parts of the request that need no seam "+ + "are not held hostage by the part that does") + + // --- AND: NOTHING was acted on. + assert.Equal(t, activitiesBefore, rowCount(t, h.db, "outbound_activities"), + "no new activity: with no seam wired there is nothing to send, and a tier that "+ + "invented one would be sending an irreversible Delete from a code path nobody "+ + "has reviewed for this deployment") + assert.Empty(t, activityOfKind(t, h.db, mtAuthorDID, "Delete"), + "no Delete{Person} in particular: peers that honour one cannot restore what they drop") + assert.Zero(t, undoActivitiesFor(t, h.db, mtAuthorDID), + "and no Undo for the vote a peer still holds") + assert.Equal(t, "delivered", voteState(t, h.db, liveVote), + "whose ledger row is UNTOUCHED: flipping it to 'undone' without sending the Undo "+ + "would hide the vote from the enumeration the real purge runs when the seam "+ + "lands — the erasure would then miss the one vote it was written for") + + assert.Equal(t, http.StatusOK, actorDocStatus(t, h, mtAuthorDID), + "and the actor document still resolves: the 410 is the destructive tier's own "+ + "statement, and making it here would tell every peer to stop resolving an "+ + "identity whose content nobody was ever asked to purge") +} + +// actorCount reports how many ap_actors rows exist for a DID. +func actorCount(t *testing.T, db *sql.DB, did string) int { + t.Helper() + var n int + require.NoError(t, db.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM ap_actors WHERE did = $1`, did).Scan(&n)) + return n +} diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go index bfee962..bbc5137 100644 --- a/internal/outbound/worker.go +++ b/internal/outbound/worker.go @@ -528,7 +528,12 @@ func (w *Worker) poison(ctx context.Context, delivery *store.OutboundDelivery, c // an outcome class, not an error class — park and parkCausal already use the // same column for held-not-failed states — and it is the durable fact that lets // a retry finish the job without repeating the POST. -const deliveryLedgerUnsettled = "ledger_unsettled" +// +// It is the STORE's constant, not a copy: every cancellation excludes rows +// carrying this class, and a second spelling here would let the resume path and +// the cancellations disagree about which rows are held — the disagreement being +// silent, and fatal in the direction where a cancel wins. +const deliveryLedgerUnsettled = store.DeliveryHeldForSettlement // statusOf is the status a held delivery was accepted with, so its settlement // records the same outcome the wire actually produced. diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go index 6aa8b1a..66cfc46 100644 --- a/internal/store/outbound_deliveries.go +++ b/internal/store/outbound_deliveries.go @@ -274,6 +274,51 @@ func (r *postgresOutboundDeliveries) CancelForActorTx(ctx context.Context, tx *s return cancelForActor(ctx, tx, actorDID) } +// DeliveryHeldForSettlement is the last_error_class of a delivery the PEER HAS +// ALREADY ACCEPTED whose local settlement — the causal stamp, the vote ledger — +// has not committed yet. The row deliberately stays `pending` so a worker can +// re-claim it and finish the bookkeeping WITHOUT repeating the POST (task 17b). +// +// The column is this package's, so the vocabulary lives here and the worker +// reads it from here: two copies of the string would let a cancellation and a +// resume disagree about which rows are held, which is exactly the bug the +// predicate below exists to prevent. +const DeliveryHeldForSettlement = "ledger_unsettled" + +// notHeldForSettlement is the term EVERY cancellation carries, and it is one +// constant rather than three because forgetting it is silent. +// +// A cancellation answers "this must not go out". A held delivery already WENT +// out: the peer holds the activity, and the only thing outstanding is our own +// record of that. Cancelling is terminal, so the worker never returns to it and +// the settlement is stranded — and the stranding is invisible until it surfaces +// somewhere else entirely. Two places, both crossing sub-run boundaries: the +// vote reseed subtracts only `delivered` rows, so a stranded one over-counts a +// served score forever; and the destructive tier enumerates a purged actor's +// live votes from that same column, so the Undo an erasure owes the peer is +// never enqueued and an erased user's vote stands on an instance nobody told. +// +// It is the same rule as "terminal rows are untouched", one state over: a +// decision that arrives LATER may not rewrite the record of something that has +// already happened. There is nothing to stop here — the send is done. +// +// IS DISTINCT FROM rather than <>, and the reason is NOT that NULLs exist: +// last_error_class is NOT NULL DEFAULT ” (migration 020), so a plain <> is +// correct today and both forms cancel the ordinary never-failed row. The +// NULL-safe form is used because this one fragment is pasted into every +// cancellation there is, and under <> the day that column becomes nullable is +// the day EVERY cancellation silently stops matching the rows it exists to +// cancel — a failure that shows up as deliveries going out after a user asked +// us to stop, nowhere near the schema change that caused it. +// +// The value is interpolated from a CONSTANT and never from input, which is what +// lets one fragment drop into statements with different parameter counts. The +// column name is unqualified deliberately: outbound_activities (the only table +// any of these statements joins) has no such column, so it is unambiguous +// everywhere and stays correct whether the target is aliased or not. +const notHeldForSettlement = ` + AND last_error_class IS DISTINCT FROM '` + DeliveryHeldForSettlement + `'` + // cancelForActor is the consent/kill-switch withdrawal: park the actor's // PENDING work as cancelled (never poisoned — this is not a failure) across // EVERY community they have work in, because the decision is about the actor. @@ -286,7 +331,7 @@ func cancelForActor(ctx context.Context, ex execer, actorDID string) (int64, err FROM outbound_activities a WHERE d.activity_id = a.activity_id AND a.actor_did = $1 - AND d.state = 'pending'` + AND d.state = 'pending'` + notHeldForSettlement result, err := ex.ExecContext(ctx, query, actorDID) if err != nil { @@ -305,7 +350,7 @@ func (r *postgresOutboundDeliveries) CancelForCommunity(ctx context.Context, ord query := ` UPDATE outbound_deliveries SET state = 'cancelled', claimed_until = NULL, updated_at = now() - WHERE ordering_key = $1 AND state = 'pending'` + WHERE ordering_key = $1 AND state = 'pending'` + notHeldForSettlement return r.cancel(ctx, "cancel outbound_deliveries for community", query, orderingKey) } @@ -336,7 +381,7 @@ func cancelPendingForActorInCommunity(ctx context.Context, ex execer, actorDID, WHERE d.activity_id = a.activity_id AND a.actor_did = $1 AND d.ordering_key = $2 - AND d.state = 'pending'`, actorDID, orderingKey) + AND d.state = 'pending'`+notHeldForSettlement, actorDID, orderingKey) if err != nil { return 0, fmt.Errorf("cancel outbound_deliveries for %q in %q: %w", actorDID, orderingKey, err) } @@ -458,7 +503,7 @@ func (r *postgresOutboundDeliveries) CancelClaimed(ctx context.Context, activity SET state = 'cancelled', claimed_until = NULL, updated_at = now() WHERE activity_id = $1 AND target_inbox = $2 AND state = 'pending' - AND claimed_until = $3 + AND claimed_until = $3` + notHeldForSettlement + ` RETURNING 1 ) SELECT diff --git a/internal/votes/optout_settlement_test.go b/internal/votes/optout_settlement_test.go new file mode 100644 index 0000000..90837d6 --- /dev/null +++ b/internal/votes/optout_settlement_test.go @@ -0,0 +1,163 @@ +package votes + +import ( + "context" + stderrors "errors" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/outbound" + "tidepool/internal/store" +) + +// TASK 17d, CYCLE 3 — THE ledger_unsettled PATH IS NOT A CONSENT QUESTION. +// +// 17b established a state no other delivery reaches: a POST the peer CONFIRMED +// whose local settlement failed. That delivery is HELD — last_error_class +// ledger_unsettled, still pending, still re-claimable — and it RESUMES AT THE +// SETTLEMENT, deliberately skipping the kill switch, the causal gate and the +// consent recheck, because all three answer "should this go out?" and the peer +// already answered it. +// +// 17d adds a consent-shaped reason to cancel work: an opt-out disables the actor +// and cancels their queued deliveries. Both halves of that must stay off this +// path, and the harm is the same in either direction — the vote is AT LEMMY, and +// a delivery that never settles leaves outbound_votes saying 'pending' forever: +// +// - 17b's reseed subtracts only DELIVERED rows, so the fediverse-only tally it +// serves keeps counting our persona's vote as a stranger's, permanently, on +// a score readers see; +// - and 17d's own destructive tier enumerates live votes from that same column, +// so the Undo that erasure owes the peer is never enqueued — the purged +// actor's vote stands on Lemmy forever, which is precisely what the tier +// exists to prevent. +// +// The two tests below are the two doors into that outcome. Each drives the vote +// through the REAL path (consumer → enqueue → worker → wire) with the ledger +// write faulted, so the held state is one the system actually produced. + +// heldForSettlement reports the single delivery's state and outcome class. +func heldForSettlement(t *testing.T, l *lifecycle) (state, class string) { + t.Helper() + require.NoError(t, l.db.QueryRow( + `SELECT state, COALESCE(last_error_class, '') FROM outbound_deliveries`).Scan(&state, &class)) + return state, class +} + +// newHeldVote casts one persona down-vote, lets it reach the peer, and faults +// the ledger write behind it. It returns the fault switch so a caller can clear +// it once the state under test has been established. +func newHeldVote(t *testing.T) (*lifecycle, *faultyVotes) { + t.Helper() + faulty := &faultyVotes{err: stderrors.New("ledger write failed")} + l := newLifecycle(t, 5, func(o *outbound.WorkerOptions) { + faulty.OutboundVotes = store.NewOutboundVotes(o.DB) + faulty.failSet = true + o.Votes = faulty + }) + ctx := context.Background() + + // A Lemmy human's live up-vote, so the subject has a real score to be wrong + // about. + require.NoError(t, l.agg.ApplyVote(ctx, like(activityID(t, 1), voterAlice, tpSubject), "")) + + l.castVote(t, "3lztprev00001", directionDown) + deliverTolerating(t, l) + + require.Equal(t, []string{"Dislike"}, l.sender.kinds(), + "precondition: the vote really did reach the peer — everything here is about what "+ + "happens AFTER the wire said yes") + state, class := heldForSettlement(t, l) + require.Equal(t, "pending", state, + "precondition: a settlement that fails after a confirmed POST HOLDS the delivery — "+ + "non-terminal, so the worker comes back to it") + require.Equal(t, "ledger_unsettled", class, + "precondition: and it carries the class that tells the next claim to resume at the "+ + "settlement rather than at the wire") + require.Equal(t, string(store.DeliveredStatePending), l.state(t), + "precondition: the ledger row is the one the hold exists to settle, and it is unsettled") + + return l, faulty +} + +// settleAndAssert clears the fault, gives the worker its next look, and asserts +// the ledger caught up. The message is the same for both doors because the +// damage is: the peer holds this vote, and only this row can say so. +func settleAndAssert(t *testing.T, l *lifecycle, faulty *faultyVotes) { + t.Helper() + faulty.failSet = false + releaseHeldDelivery(t, l) + deliverTolerating(t, l) + + assert.Equal(t, string(store.DeliveredStateDelivered), l.state(t), + "the held delivery must still settle: Lemmy holds this vote, and outbound_votes is "+ + "the only record of that. A row stranded at 'pending' over-counts the served score "+ + "forever AND hides the vote from the destructive tier's Undo enumeration, so an "+ + "erased user's vote stands on a peer that was never told") + assert.Equal(t, []string{"Dislike"}, l.sender.kinds(), + "and it settles WITHOUT going back to the wire: the peer accepted this activity once, "+ + "and the resume path exists because re-POSTing it is not the recovery") +} + +// TestAHeldSettlementResumesEvenForADisabledActor is the WORKER door. +// +// The actor is disabled and opted out — everything the claim-time consent +// recheck reads says "do not speak for this user" — and that recheck must not be +// reachable from the resume path. A consent check re-introduced ahead of it +// cancels the exact row the delivery was held to settle. +func TestAHeldSettlementResumesEvenForADisabledActor(t *testing.T) { + l, faulty := newHeldVote(t) + ctx := context.Background() + + // The two facts the consent recheck reads, written directly: this test is + // about the WORKER's decision, so the opt-out arrives here as state rather + // than as an event with side effects of its own (that is the next test). + require.NoError(t, store.NewAPActors(l.db).SetEnabled(ctx, tpNativeDID, false)) + _, err := store.NewFederationPrefs(l.db).Upsert(ctx, store.FederationPref{ + DID: tpNativeDID, Enabled: false, Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + + settleAndAssert(t, l, faulty) + + state, _ := heldForSettlement(t, l) + assert.Equal(t, "delivered", state, + "and the delivery reaches its terminal state rather than being cancelled: cancelling "+ + "a delivery the peer already accepted is not a withdrawal, it is a lost record of "+ + "something that happened") +} + +// TestAnOptOutDoesNotCancelADeliveryHeldForSettlement is the CONSUMER door, and +// the one 17d opened. +// +// handleFederation cancels this actor's PENDING deliveries, and a held +// settlement is pending by construction — that is what makes it re-claimable. +// Cancelling is terminal, so the worker never comes back, and the ledger row the +// hold existed to settle is stranded at 'pending' forever. +// +// The cancellation is right about everything it was written for: work that has +// NOT gone out must not go out. This row already went out. "Stop sending" and +// "forget what was sent" are different instructions, and only the first one was +// asked for. +func TestAnOptOutDoesNotCancelADeliveryHeldForSettlement(t *testing.T) { + l, faulty := newHeldVote(t) + + // The user opts out, through the real consumer: the soft tier, which + // disables the actor and cancels their queued work in one transaction. + l.handle(t, fmt.Sprintf( + `{"did":%q,"time_us":9700,"kind":"commit","commit":{"rev":"3lztprev00009","operation":"create",`+ + `"collection":"social.coves.bridge.federation","rkey":"self","cid":%q,`+ + `"record":{"$type":"social.coves.bridge.federation","enabled":false}}}`, + tpNativeDID, testCID)) + + enabled, err := store.NewAPActors(l.db).GetByDID(context.Background(), tpNativeDID) + require.NoError(t, err) + require.False(t, enabled.Enabled, + "precondition: the opt-out really was applied — the assertions below must not pass "+ + "because nothing happened") + + settleAndAssert(t, l, faulty) +}