From 97c413b480c3c70c773ab4c7c26636ce3dd216e9 Mon Sep 17 00:00:00 2001 From: Bretton Date: Fri, 14 Aug 2026 17:52:06 -0700 Subject: [PATCH] fix(optout): a withdrawal that survives its own replay, and stays withdrawn (17d review) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The destructive tier worked on the happy path and failed in three ways that only appear once something goes wrong. All three came out of review. THE PURGE WAS NOT REPLAY-SAFE, and the cancel-before-purge ordering adopted after the deadlock is what broke it. The purge commits on its own transaction and the gate commits later, so a failure between them replays the record. On replay the sweeping cancel ran FIRST and cancelled the Delete{Person} fan-out and the vote Undos the previous purge had already committed; EnqueueTx then refused to revive them by design, and the enumeration returned nothing because those votes were already marked undone. Net result: actor tombstoned, prefs recorded, ZERO activities ever delivered, logged as "destructive opt-out applied". The same missing predicate over-reached in the soft tier, where opting out cancelled a user's own pending self-delete and left the post standing on Lemmy forever. The cancel is now two statements named for two different decisions. CancelForActor stays sweeping for the kill switch and the operator's manual cancel, which mean "stop the queue". CancelOutwardForActorTx is consent: it cancels what PUBLISHES and lets store.RetractionKinds go out. worker.isRetraction reads that same list, so the queue and the worker cannot drift about what a stopped user is still owed. That predicate also dissolves the ordering constraint rather than relocating it — verified the hard way, since biting it does not fail the replay test, it DEADLOCKS it, which is the deadlock the ordering existed to avoid. A VOTE HELD FOR SETTLEMENT WAS INVISIBLE TO THE ERASURE. A hold is a delivery the peer ACCEPTED — only our bookkeeping is outstanding — but the enumeration read delivered_state alone, so the purge owed nothing while the peer held the Like. ListStandingForActor now reads delivered UNION held-for-settlement. AND UNDONE WAS NOT ACTUALLY TERMINAL. The enumeration fix alone is undone by the very settlement it races: the held delivery resumes after the withdrawal and writes 'delivered' over the retraction, so 17b's reseed keeps subtracting a vote from a served score on behalf of somebody who no longer exists. The guard now lives in SetDeliveredState, and its two zero-row cases are kept apart deliberately — an already-retracted row is a DECIDED no-op and reports success, because an error there holds the delivery forever retrying a write that can never apply, while a genuinely missing row still reports NotFound. The 410 was observable before the withdrawal was delivered: a peer that re-dereferenced the actor to verify the signature on the Delete got Gone for the message announcing that deletion. Both actor routes and the outbox now answer with an AS2 Tombstone carrying formerType, deleted, and the public key — verification material, none of the profile. Keeping the full document served until every delivery is terminal was considered and rejected: it publishes an erased user's document for as long as any peer is down, which is a worse failure for an erasure tier than the one it fixes. The residual is in FOLLOWUPS. Also: a tombstoned actor can no longer be re-enabled (forced in the statement, so no caller gets a spurious NotFound) and a purged preference cannot be deleted; federation_prefs.purged_at separates requested from committed and is stamped inside the purge's own transaction; the confirmation resolver requires https, refuses userinfo/query/fragment, matches the #atproto_pds fragment exactly, and pins every redirect hop to the named authority, with every violation an error rather than a verdict; one unresolvable community no longer holds an entire erasure hostage; and the purge resolves every inbox BEFORE opening its transaction, which is why the ingest package took 332s and failed while each test passed alone. Tests: whole-package -count=2 across store/outbound/consume/ingest, full suite 21/21, e2e 284s. Every guard tooth-checked, including one bite of mine that was itself invalid (an untyped parameter made Postgres fail type inference and reddened all five tests for a reason that had nothing to do with the guard). Co-Authored-By: Claude Opus 5 (1M context) --- FOLLOWUPS.md | 35 ++ internal/consume/account_confirm_test.go | 20 +- internal/consume/account_status.go | 98 +++- .../consume/account_status_transport_test.go | 481 ++++++++++++++++++ internal/consume/federation.go | 53 +- internal/consume/rev_gate.go | 18 +- .../migrations/030_federation_pref_purged.sql | 32 ++ internal/ingest/optout_retraction_test.go | 181 +++++++ internal/optout/terminator.go | 35 +- internal/outbound/enqueuer.go | 23 +- internal/outbound/fanout_test.go | 10 +- internal/outbound/purge.go | 129 ++++- internal/outbound/purge_held_vote_test.go | 224 ++++++++ internal/outbound/worker.go | 8 +- internal/personas/serving.go | 51 +- internal/store/ap_actors.go | 20 +- internal/store/federation_prefs.go | 84 ++- internal/store/interfaces.go | 63 ++- internal/store/models.go | 31 +- internal/store/outbound_deliveries.go | 81 ++- internal/store/outbound_standing_row_test.go | 237 +++++++++ internal/store/outbound_vote_terminal_test.go | 127 +++++ internal/store/outbound_votes.go | 75 ++- 23 files changed, 2001 insertions(+), 115 deletions(-) create mode 100644 internal/consume/account_status_transport_test.go create mode 100644 internal/db/migrations/030_federation_pref_purged.sql create mode 100644 internal/ingest/optout_retraction_test.go create mode 100644 internal/outbound/purge_held_vote_test.go create mode 100644 internal/store/outbound_standing_row_test.go create mode 100644 internal/store/outbound_vote_terminal_test.go diff --git a/FOLLOWUPS.md b/FOLLOWUPS.md index ddcecbe..cb7d58e 100644 --- a/FOLLOWUPS.md +++ b/FOLLOWUPS.md @@ -369,3 +369,38 @@ task documents and git history rather than this list. rejected as the design: its op list would not match the empty MST diff, which sync-v1.1-validating relays may refuse). Remaining follow-up: the e2e harness has no scenario covering it. + +## Opt-out lifecycle (17d) + +- **The withdrawal's 410 is observable before the withdrawal is delivered.** + The destructive tier commits the actor tombstone together with the + `Delete{Person}` enqueues, but the worker POSTs minutes later. A peer that + re-dereferences the actor to verify the signature on that Delete therefore + gets `410 Gone` — for the message announcing the very deletion it is trying + to authenticate. The tier can revoke its own precondition. + + MITIGATED, not closed: both actor routes (and `handleOutbox`, which would + otherwise invite the re-fetch loop the 410 exists to end) now answer 410 + with an AS2 `Tombstone` carrying `formerType: Person`, `deleted`, and **the + public key** — so a peer that reads the body can still verify. + + RULED against the "complete" fix of keeping the actor document served until + every person-delete delivery is terminal. That keeps an erased user's + document published for as long as any peer is down — potentially forever — + which is a worse failure for an erasure tier than the one it fixes. The + Tombstone hands over verification material and none of the profile, which is + the better trade, and it is why the committed test pinning "410 immediately, + with no worker run" was deliberately left standing. + + RESIDUAL, and it is genuinely not ours to close: a peer that branches on the + status code alone, without reading the body, still cannot verify the + withdrawal. Revisit only if a real implementation (Lemmy specifically) is + observed dropping our person-deletes for this reason — the fix would then be + peer-shaped (retry after the delete, cached-key acceptance), not a change to + when we serve the tombstone. + +- **A delivery claimed and mid-POST at purge time is indistinguishable from + one that will never be sent.** The purge enumerates what the peer holds + (delivered, plus held-for-settlement), but a delivery the worker has claimed + and is POSTing right now is neither. Detecting it needs the peer's state, + not ours. Belongs to 17e (reconciliation). diff --git a/internal/consume/account_confirm_test.go b/internal/consume/account_confirm_test.go index ba0b557..6852a14 100644 --- a/internal/consume/account_confirm_test.go +++ b/internal/consume/account_confirm_test.go @@ -12,6 +12,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" + "tidepool/internal/errors" "tidepool/internal/optout" "tidepool/internal/store" ) @@ -120,12 +121,19 @@ func appliedSeq(t *testing.T, database *sql.DB, did string) int64 { } // storedPref reads the terminal preference, or nil when none was recorded. +// +// NOT-FOUND IS THE ONLY ABSENCE. Swallowing every error here would make +// "nothing was recorded" true of a query that failed, a table that was dropped +// and a typo in the column list — three different nothings, all reading as the +// assertion passing, on the tests whose whole content is that a row does not +// exist. 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 { + if errors.IsNotFound(err) { return nil } + require.NoError(t, err, "read the federation preference for %s", did) return pref } @@ -167,10 +175,10 @@ func TestConfirmedDeletion_RunsTheDestructivePathAndClosesTheEvent(t *testing.T) } // --------------------------------------------------------------------------- -// (2) CONFIRMED STILL LIVE — nothing destructive, and the seq STILL ADVANCES +// (2) CONFIRMED LIVE — nothing destructive, and the seq STILL ADVANCES // --------------------------------------------------------------------------- -func TestUnconfirmedDeletionOfALiveAccountDoesNothingAndStillAdvances(t *testing.T) { +func TestConfirmedLiveAccountIsLeftAloneAndTheEventStillAdvances(t *testing.T) { database := dispatchTestDB(t) seedAPActor(t, database, dispatchNativeDID, "alice") world := newTerminalWorld(t, database, true) @@ -204,13 +212,13 @@ func TestUnconfirmedDeletionOfALiveAccountDoesNothingAndStillAdvances(t *testing } // --------------------------------------------------------------------------- -// (3) CONFIRM FAILED — retryable, nothing sent, seq NOT advanced +// (3) UNKNOWN (the confirm failed) — retryable, nothing sent, seq NOT advanced // --------------------------------------------------------------------------- -// TestUnconfirmableDeletionIsRetryableAndAdvancesNothing is the case the seam +// TestUnknownConfirmOutcomeIsRetryableAndAdvancesNothing 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) { +func TestUnknownConfirmOutcomeIsRetryableAndAdvancesNothing(t *testing.T) { database := dispatchTestDB(t) seedAPActor(t, database, dispatchNativeDID, "alice") world := newTerminalWorld(t, database, true) diff --git a/internal/consume/account_status.go b/internal/consume/account_status.go index 65168da..7bd2a43 100644 --- a/internal/consume/account_status.go +++ b/internal/consume/account_status.go @@ -21,7 +21,7 @@ import ( // atprotoPDSServiceID is the DID document service entry that names a repo's // hosting PDS. -const atprotoPDSServiceID = "#atproto_pds" +const atprotoPDSServiceID = "atproto_pds" // maxRepoStatusBytes caps the PDS response. getRepoStatus answers a few dozen // bytes; the host is named by a document a stranger controls, so it is not read @@ -82,21 +82,62 @@ type didService struct { } // pdsEndpoint finds the repo's hosting PDS. The id suffix is what identifies it -// (documents spell it "#atproto_pds" or the full "did:plc:xxx#atproto_pds"), -// and the endpoint must be an absolute http(s) URL before it reaches a request: -// it comes from a document its own subject controls. -func pdsEndpoint(document *didDocument) (string, error) { +// (documents spell it "#atproto_pds" or the full "did:plc:xxx#atproto_pds"). +// +// THE ENDPOINT IS UNTRUSTED TEXT IN THE ORDINARY CASE, not merely under attack: +// it is written by the DID's own controller, which is the party a deletion +// verdict is about. So it is validated to exhaustion BEFORE the first packet — +// a guard that fires on the RESPONSE has already sent the bridge somewhere a +// stranger chose, and has already leaked which DID it is about to act on. +func pdsEndpoint(document *didDocument) (*url.URL, error) { for _, service := range document.Service { - if !strings.HasSuffix(service.ID, atprotoPDSServiceID) { + // EXACT fragment match, not a suffix test: "#not_atproto_pds" ends with + // the same characters, and a document its own subject writes is where + // that shows up. Both spellings are legal — the bare "#atproto_pds" and + // the fully-qualified "did:plc:xxx#atproto_pds" — so what is compared is + // the fragment itself. + if _, fragment, found := strings.Cut(service.ID, "#"); !found || fragment != atprotoPDSServiceID { continue } - parsed, err := url.Parse(service.ServiceEndpoint) - if err != nil || parsed.Host == "" || (parsed.Scheme != "http" && parsed.Scheme != "https") { - return "", fmt.Errorf("DID document names a PDS endpoint that is not an absolute http(s) URL") - } - return strings.TrimSuffix(service.ServiceEndpoint, "/"), nil + return checkedPDSEndpoint(service.ServiceEndpoint) + } + return nil, fmt.Errorf("DID document names no atproto PDS") +} + +// checkedPDSEndpoint accepts only a plain, absolute HTTPS origin. +// +// HTTPS IS NOT A PREFERENCE HERE. This response decides whether a user's content +// is erased from every instance that holds it, and no peer un-deletes. Over +// cleartext, anyone on the path can WRITE {"active":false,"status":"deleted"} — +// the irreversible verdict, for free, with no credential and no compromise of +// either endpoint. There is no allow-insecure option on purpose: a deployment +// whose PDS is reachable only over http cannot confirm deletions, and it fails +// CLOSED (a retryable error an operator can see) rather than acting on an answer +// nobody can vouch for. +// +// Userinfo is refused too: credentials in a URL a stranger wrote are not ours to +// send, and Go would put them on the wire. +func checkedPDSEndpoint(endpoint string) (*url.URL, error) { + parsed, err := url.Parse(endpoint) + switch { + case err != nil: + return nil, fmt.Errorf("DID document names an unparseable PDS endpoint") + case parsed.Scheme != "https": + return nil, fmt.Errorf( + "DID document names a PDS endpoint that is not https; a confirmation over a channel "+ + "anyone can rewrite cannot decide an irreversible erasure (%q)", parsed.Scheme) + case parsed.Host == "": + return nil, fmt.Errorf("DID document names a PDS endpoint with no host") + case parsed.User != nil: + return nil, fmt.Errorf("DID document names a PDS endpoint carrying userinfo") + case parsed.RawQuery != "" || parsed.Fragment != "": + // An XRPC base is an origin with an optional path prefix. A query or + // fragment on it means the value is not that, and appending to it would + // silently produce a URL nobody wrote. + return nil, fmt.Errorf("DID document names a PDS endpoint carrying a query or fragment") } - return "", fmt.Errorf("DID document names no atproto PDS") + parsed.Path = strings.TrimSuffix(parsed.Path, "/") + return parsed, nil } // repoStatus is the sliver of com.atproto.sync.getRepoStatus this reads. @@ -117,9 +158,14 @@ const repoStatusDeleted = "deleted" // 400s an unrecognised repo looks identical to one that is misconfigured, and // "the host said something we did not understand" is not evidence a user // deleted their account. -func (r *HandleResolver) repoDeleted(ctx context.Context, endpoint, did string) (bool, error) { - target := endpoint + "/xrpc/com.atproto.sync.getRepoStatus?did=" + url.QueryEscape(did) - request, err := http.NewRequestWithContext(ctx, http.MethodGet, target, nil) +func (r *HandleResolver) repoDeleted(ctx context.Context, endpoint *url.URL, did string) (bool, error) { + // BUILT FROM THE PARSED URL, never by concatenation: JoinPath escapes the + // segments and preserves a path prefix (https://host/pds), and setting the + // query as a value keeps the did from being pasted into a string that may + // already contain one. + target := endpoint.JoinPath("/xrpc/com.atproto.sync.getRepoStatus") + target.RawQuery = url.Values{"did": []string{did}}.Encode() + request, err := http.NewRequestWithContext(ctx, http.MethodGet, target.String(), nil) if err != nil { return false, fmt.Errorf("build repo status request for %s: %w", did, err) } @@ -128,7 +174,27 @@ func (r *HandleResolver) repoDeleted(ctx context.Context, endpoint, did string) request.Header.Set("User-Agent", r.userAgent) } - response, err := r.httpClient.Do(request) + // EVERY HOP STAYS ON THE NAMED PDS, over https. The DID document naming that + // PDS is the ENTIRE authorization for this answer — it is what makes the + // response evidence about this repo rather than an opinion from a stranger — + // and a redirect is written by the very server whose answer we are trying to + // verify, so "it told us to" is exactly as trustworthy as the answer itself. + // + // The client is COPIED rather than mutated: this policy belongs to this + // request, and the same client also fetches DID documents and well-knowns, + // where redirects are ordinary. Copying shares the transport (and its SSRF + // guard) while giving this call its own rules. + client := *r.httpClient + client.CheckRedirect = func(hop *http.Request, _ []*http.Request) error { + if hop.URL.Scheme != "https" || !strings.EqualFold(hop.URL.Host, request.URL.Host) { + return fmt.Errorf( + "refusing redirect to %s: a deletion verdict may come only from the PDS the DID "+ + "document names, over https", hop.URL.Redacted()) + } + return nil + } + + response, err := client.Do(request) if err != nil { return false, fmt.Errorf("fetch repo status for %s: %w", did, err) } diff --git a/internal/consume/account_status_transport_test.go b/internal/consume/account_status_transport_test.go new file mode 100644 index 0000000..4227205 --- /dev/null +++ b/internal/consume/account_status_transport_test.go @@ -0,0 +1,481 @@ +package consume + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/ap" +) + +// TASK 17d — THE CONFIRM AT THE TRANSPORT (decision 19). +// +// The four dispatcher-tier cases in account_confirm_test.go pin what the bridge +// DOES with each of the confirm's three answers. They cannot pin where those +// answers come from: the confirmer is stubbed there, so the suite stays green +// against a resolver that picks the wrong service entry, follows a redirect off +// the PDS, accepts plain HTTP, or turns a non-200 into a verdict. +// +// THIS IS THE LAYER WHERE THE DANGER LIVES, and the danger is specific: the seam +// was built so that no FAILURE could become a VERDICT. Every route by which +// somebody other than the DID's own PDS gets to answer hands that verdict away. +// `{"active":false,"status":"deleted"}` from an on-path attacker, or from a +// server reached after an HTTPS→HTTP downgrade, is the irreversible verdict — +// Delete{Person, removeData:true} to every instance holding that user's content, +// which no peer undoes. +// +// So the rule these tests describe is one rule: THE ANSWER MUST COME FROM THE +// PDS THE DID DOCUMENT NAMES, OVER A CHANNEL THAT CANNOT BE READ OR REWRITTEN, +// AND ANYTHING ELSE IS AN ERROR. Not a "live" verdict — an error, because "we +// could not confirm" is its own outcome. +// +// Local-only: one httptest listener plays the PLC directory and every PDS, and +// the transport REFUSES any host the fixture does not serve, so a policy bug +// cannot become a request to the real internet. + +const ( + asDID = "did:plc:accountstatus0000001" + // The PDS the document names, and a second one that has no claim on this + // DID. A redirect from the first to the second is the on-path case in its + // most honest form: the peer that answers is not the peer we asked. + asPDSHost = "pds.example" + asOtherPDSHost = "other-pds.example" + + asStatusPath = "/xrpc/com.atproto.sync.getRepoStatus" +) + +// asAttempt is one outbound request the resolver made, recorded BEFORE the test +// transport rewrites it onto the local listener — so the scheme and host are the +// ones the resolver actually asked for. +type asAttempt struct { + scheme string + host string + path string +} + +// pdsWorld is a PLC directory plus one or more PDS hosts. +type pdsWorld struct { + server *httptest.Server + + mu sync.Mutex + // The DID document: what endpoint it names, or a raw body / status override. + endpoint string + omitPDS bool + plcStatus int + plcBody string + // Per-PDS-host behaviour for getRepoStatus. + status map[string]int + body map[string]string + redirect map[string]string + + attempts []asAttempt +} + +func newPDSWorld(t *testing.T) *pdsWorld { + t.Helper() + world := &pdsWorld{ + endpoint: "https://" + asPDSHost, + status: map[string]int{}, + body: map[string]string{}, + redirect: map[string]string{}, + } + + mux := http.NewServeMux() + mux.HandleFunc(asStatusPath, func(w http.ResponseWriter, r *http.Request) { + host := hostOnly(r.Host) + world.mu.Lock() + location, redirects := world.redirect[host] + status, forcedStatus := world.status[host] + body, forcedBody := world.body[host] + world.mu.Unlock() + + if redirects { + http.Redirect(w, r, location, http.StatusFound) + return + } + w.Header().Set("Content-Type", "application/json") + if forcedStatus { + w.WriteHeader(status) + } + if forcedBody { + _, _ = w.Write([]byte(body)) + return + } + // The default answer is a LIVE repo, so a test that reaches the wrong + // server by accident cannot pass by inheriting a deleted verdict. + _, _ = w.Write([]byte(`{"did":"` + asDID + `","active":true}`)) + }) + + // Everything else is the PLC directory: GET /{did}. + mux.HandleFunc("/", func(w http.ResponseWriter, r *http.Request) { + world.mu.Lock() + status, endpoint, omit, body := world.plcStatus, world.endpoint, world.omitPDS, world.plcBody + world.mu.Unlock() + + if status != 0 { + w.WriteHeader(status) + _, _ = w.Write([]byte(`{"message":"forced directory failure"}`)) + return + } + w.Header().Set("Content-Type", "application/did+ld+json") + if body != "" { + _, _ = w.Write([]byte(body)) + return + } + service := fmt.Sprintf( + `{"id":"#atproto_pds","type":"AtprotoPersonalDataServer","serviceEndpoint":%q}`, endpoint) + if omit { + // A document whose only service is something else entirely. It is + // NOT evidence of a deletion — the same shape appears while a + // document is mid-rewrite or hosting is moving. + service = `{"id":"#atproto_labeler","type":"AtprotoLabeler","serviceEndpoint":"https://labeler.example"}` + } + _, _ = w.Write([]byte(fmt.Sprintf( + `{"id":%q,"alsoKnownAs":["at://alice.example"],"service":[%s]}`, asDID, service))) + }) + + world.server = httptest.NewServer(mux) + t.Cleanup(world.server.Close) + return world +} + +// resolver builds the REAL HandleResolver over this world. +func (w *pdsWorld) resolver(t *testing.T) *HandleResolver { + t.Helper() + resolver, err := NewHandleResolver(ResolverOptions{ + PLCDirectoryURL: w.server.URL, + HTTPClient: &http.Client{Transport: &asRecordingTransport{world: w}}, + UserAgent: "tidepool-test/0.1", + LookupTXT: func(context.Context, string) ([]string, error) { return nil, nil }, + }) + require.NoError(t, err) + return resolver +} + +func (w *pdsWorld) record(a asAttempt) { + w.mu.Lock() + defer w.mu.Unlock() + w.attempts = append(w.attempts, a) +} + +// statusAttempts lists the getRepoStatus requests the resolver made. The DID +// document fetch is excluded: every case here makes exactly one of those, and +// what is under test is who got asked for the VERDICT. +func (w *pdsWorld) statusAttempts() []asAttempt { + w.mu.Lock() + defer w.mu.Unlock() + var out []asAttempt + for _, attempt := range w.attempts { + if strings.HasPrefix(attempt.path, asStatusPath) { + out = append(out, attempt) + } + } + return out +} + +// asRecordingTransport records what the resolver asked for and then serves it +// from the local listener. It refuses every host outside the fixture, so the +// tests stay offline even while proving a guard is missing. +type asRecordingTransport struct{ world *pdsWorld } + +func (rt *asRecordingTransport) RoundTrip(req *http.Request) (*http.Response, error) { + host := hostOnly(req.URL.Host) + rt.world.record(asAttempt{scheme: req.URL.Scheme, host: host, path: req.URL.Path}) + + listener := hostOnly(rt.world.server.Listener.Addr().String()) + if host != listener && !strings.HasSuffix(host, ".example") { + return nil, fmt.Errorf("refusing outbound request to %s: these tests never touch a network", req.URL) + } + clone := req.Clone(req.Context()) + clone.Host = req.URL.Host // the mux dispatches PDS hosts by Host header + clone.URL.Scheme = "http" + clone.URL.Host = rt.world.server.Listener.Addr().String() + return http.DefaultTransport.RoundTrip(clone) +} + +// --------------------------------------------------------------------------- +// The verdicts themselves +// --------------------------------------------------------------------------- + +// TestAccountStatus_OnlyAnInactiveDeletedRepoIsDeleted pins the predicate both +// halves of it: `active` alone reads a suspension as a deletion, and `status` +// alone believes a field a PDS may omit for a live repo. +func TestAccountStatus_OnlyAnInactiveDeletedRepoIsDeleted(t *testing.T) { + for _, tc := range []struct { + name string + body string + want bool + }{ + {"live repo", `{"active":true}`, false}, + {"live repo with a status field", `{"active":true,"status":"active"}`, false}, + {"deactivated", `{"active":false,"status":"deactivated"}`, false}, + {"suspended", `{"active":false,"status":"suspended"}`, false}, + {"takendown", `{"active":false,"status":"takendown"}`, false}, + {"throttled", `{"active":false,"status":"throttled"}`, false}, + {"inactive with no status at all", `{"active":false}`, false}, + {"active AND deleted, which is incoherent", `{"active":true,"status":"deleted"}`, false}, + {"the one shape that means gone", `{"active":false,"status":"deleted"}`, true}, + } { + t.Run(tc.name, func(t *testing.T) { + world := newPDSWorld(t) + world.body[asPDSHost] = tc.body + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + require.NoError(t, err) + assert.Equal(t, tc.want, deleted, + "only active=false AND status=deleted is a deletion: every other inactive "+ + "state is one a user comes back from, and withdrawing their identity "+ + "over a suspension destroys it irreversibly") + }) + } +} + +// --------------------------------------------------------------------------- +// The directory leg: no failure may become a verdict +// --------------------------------------------------------------------------- + +func TestAccountStatus_DirectoryFailuresAreErrorsNeverVerdicts(t *testing.T) { + for _, tc := range []struct { + name string + setup func(*pdsWorld) + }{ + {"directory 503", func(w *pdsWorld) { w.plcStatus = http.StatusServiceUnavailable }}, + {"directory 404", func(w *pdsWorld) { w.plcStatus = http.StatusNotFound }}, + {"directory serves a non-document", func(w *pdsWorld) { w.plcBody = `maintenance` }}, + {"document names no atproto PDS", func(w *pdsWorld) { w.omitPDS = true }}, + } { + t.Run(tc.name, func(t *testing.T) { + world := newPDSWorld(t) + tc.setup(world) + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + require.Error(t, err, + "an unanswerable directory is UNKNOWN: read as 'live' it silently drops a "+ + "real deletion and leaves the user federated forever, read as 'deleted' "+ + "it erases someone who never left. Only an error keeps the event redrivable") + assert.False(t, deleted, "and no verdict may ride out alongside the error") + assert.Empty(t, world.statusAttempts(), + "nor may a missing or unreadable document be papered over by asking somebody "+ + "else: with no PDS named, there is no authority to ask") + }) + } +} + +// --------------------------------------------------------------------------- +// The endpoint the document names: it is attacker-influenced input +// --------------------------------------------------------------------------- + +// TestAccountStatus_AnUnusablePDSEndpointIsRefusedBeforeAnyRequest covers the +// shapes that are not an absolute http(s) URL. The serviceEndpoint is written by +// the DID's own controller, so it reaches this code as untrusted text, and the +// check has to hold BEFORE the first packet. +func TestAccountStatus_AnUnusablePDSEndpointIsRefusedBeforeAnyRequest(t *testing.T) { + for _, endpoint := range []string{ + "pds.example", // no scheme: url.Parse gives it no host + "/xrpc", // relative + "file:///etc/passwd", // not http(s) + "ftp://pds.example", // not http(s) + "", // present but empty + "https://user:pw@pds.example", // userinfo: the credential is not ours to send + } { + t.Run(endpoint, func(t *testing.T) { + world := newPDSWorld(t) + world.endpoint = endpoint + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + require.Error(t, err, "an endpoint that is not a plain absolute http(s) URL is unusable") + assert.False(t, deleted) + assert.Empty(t, world.statusAttempts(), + "and it is refused BEFORE any request: a guard that fires on the response has "+ + "already sent the bridge somewhere a stranger chose") + }) + } +} + +// TestAccountStatus_APlainHTTPPDSEndpointIsRefusedBeforeAnyRequest is the first +// of the three transport rules, and the simplest. +// +// A cleartext confirmation is a confirmation anyone on the path can WRITE. The +// answer decides whether a user's content is erased everywhere, so an attacker +// who can inject one response gets the irreversible verdict for free — no +// credentials, no compromise of either endpoint. +func TestAccountStatus_APlainHTTPPDSEndpointIsRefusedBeforeAnyRequest(t *testing.T) { + world := newPDSWorld(t) + world.endpoint = "http://" + asPDSHost + world.body[asPDSHost] = `{"active":false,"status":"deleted"}` + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + + assert.False(t, deleted, + "a cleartext answer must never produce the deleted verdict: whoever is on the path "+ + "writes it, and Delete{Person, removeData:true} is not recallable") + require.Error(t, err, + "and it is an ERROR, not a quiet 'live': the confirmation did not happen, so the "+ + "event must come back rather than be consumed") + assert.Empty(t, world.statusAttempts(), + "refused before the first packet — a request already sent has already leaked which "+ + "DID this bridge is about to act on") +} + +// TestAccountStatus_AnHTTPSToHTTPRedirectCannotProduceAVerdict is the same rule +// one hop later, and the hop is where it is usually lost: the endpoint is +// https, the policy looks satisfied, and the redirect hands the conversation to +// cleartext anyway. +func TestAccountStatus_AnHTTPSToHTTPRedirectCannotProduceAVerdict(t *testing.T) { + world := newPDSWorld(t) + world.endpoint = "https://" + asPDSHost + world.redirect[asPDSHost] = "http://" + asPDSHost + asStatusPath + "?did=" + asDID + world.body[asPDSHost] = `{"active":false,"status":"deleted"}` + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + + assert.False(t, deleted, + "a downgrade must not be able to answer: the guarantee an https endpoint bought is "+ + "gone the moment the client follows a redirect out of it, and what it bought was "+ + "the only thing standing between an on-path attacker and an irreversible erasure") + require.Error(t, err, "and the failed confirmation stays an error") + + for _, attempt := range world.statusAttempts() { + assert.Equal(t, "https", attempt.scheme, + "no request in this confirmation may go out over cleartext, redirect or not") + } +} + +// TestAccountStatus_ARedirectOffThePDSAuthorityCannotProduceAVerdict closes the +// last route: the scheme stays https and the AUTHORITY changes. +// +// The DID document names one PDS. That naming is the entire authorization for +// the answer — it is what makes the response evidence about THIS repo rather +// than an opinion from a stranger. A redirect that leaves it means the bridge +// asked the peer the user chose and believed a peer somebody else chose. +func TestAccountStatus_ARedirectOffThePDSAuthorityCannotProduceAVerdict(t *testing.T) { + world := newPDSWorld(t) + world.endpoint = "https://" + asPDSHost + world.redirect[asPDSHost] = "https://" + asOtherPDSHost + asStatusPath + "?did=" + asDID + // The host the redirect points at is delighted to confirm the deletion. + world.body[asOtherPDSHost] = `{"active":false,"status":"deleted"}` + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + + assert.False(t, deleted, + "a host the DID document never named must not be able to condemn the repo: the "+ + "redirect is written by the very server we are asking, so 'it told us to' is "+ + "exactly as trustworthy as the answer we are trying to verify") + require.Error(t, err, "an answer from the wrong authority is a failed confirmation") + + for _, attempt := range world.statusAttempts() { + assert.Equal(t, asPDSHost, attempt.host, + "and the confirmation never left the named PDS's authority at all") + } +} + +// --------------------------------------------------------------------------- +// The PDS's answer: only something we fully read and understood decides +// --------------------------------------------------------------------------- + +func TestAccountStatus_UnreadablePDSAnswersAreErrorsNeverVerdicts(t *testing.T) { + for _, tc := range []struct { + name string + status int + body string + }{ + {"500", http.StatusInternalServerError, `{"active":false,"status":"deleted"}`}, + {"400 on an unrecognised repo", http.StatusBadRequest, `{"error":"RepoNotFound"}`}, + {"404", http.StatusNotFound, ``}, + {"empty body", http.StatusOK, ``}, + {"not JSON", http.StatusOK, `hello`}, + {"JSON that is not an object", http.StatusOK, `["deleted"]`}, + { + "a body past the read cap", + http.StatusOK, + `{"padding":"` + strings.Repeat("x", maxRepoStatusBytes) + `","active":false,"status":"deleted"}`, + }, + } { + t.Run(tc.name, func(t *testing.T) { + world := newPDSWorld(t) + if tc.status != http.StatusOK { + world.status[asPDSHost] = tc.status + } + world.body[asPDSHost] = tc.body + + deleted, err := world.resolver(t).AccountStatus(context.Background(), asDID) + require.Error(t, err, + "an answer we could not read in full is not evidence of anything: a "+ + "misconfigured PDS, a proxy error page and a truncated body all look "+ + "alike, and none of them is a user deleting their account") + assert.False(t, deleted) + }) + } +} + +// --------------------------------------------------------------------------- +// The egress guard, which is the WIRING half of the same rule +// --------------------------------------------------------------------------- + +// TestAccountStatus_APrivatePDSEndpointIsRefusedByTheProductionEgressGuard +// pins the premise the tests above are allowed to relax. +// +// Everything else in this file rewrites the fixture's hosts onto a loopback +// listener, which is exactly what production must NOT do: the endpoint comes +// from a document a stranger controls, so a `serviceEndpoint` of +// http://169.254.169.254 or http://127.0.0.1:5432 would otherwise turn this +// confirmation into an SSRF probe of the bridge's own network — with the reply +// body deciding whether a user is erased. +// +// So this one wires the client production wires, and only the PDS leg goes +// through it: the directory is served locally so that a refusal here is +// attributable to the ENDPOINT rather than to the fixture being local. +func TestAccountStatus_APrivatePDSEndpointIsRefusedByTheProductionEgressGuard(t *testing.T) { + world := newPDSWorld(t) + // The endpoint the "DID document" names is the loopback listener itself, + // which is the shape of every SSRF pivot — named over HTTPS, so that the + // refusal is attributable to the ADDRESS and not to a scheme check one layer + // up. A loopback endpoint spelled http:// would be refused for the wrong + // reason and this test would advertise coverage it does not have. + world.endpoint = "https://" + world.server.Listener.Addr().String() + world.body[hostOnly(world.server.Listener.Addr().String())] = `{"active":false,"status":"deleted"}` + + resolver, err := NewHandleResolver(ResolverOptions{ + PLCDirectoryURL: world.server.URL, + HTTPClient: &http.Client{Transport: &asSplitTransport{ + directoryHost: hostOnly(world.server.Listener.Addr().String()), + directory: &asRecordingTransport{world: world}, + // The production egress: ap.NewGuardedHTTPClient(false, …), whose + // transport refuses loopback and RFC1918 at DIAL time. + guarded: ap.NewGuardedHTTPClient(false, 5*time.Second).Transport, + }}, + UserAgent: "tidepool-test/0.1", + LookupTXT: func(context.Context, string) ([]string, error) { return nil, nil }, + }) + require.NoError(t, err) + + deleted, err := resolver.AccountStatus(context.Background(), asDID) + require.Error(t, err, + "a private or loopback PDS endpoint must not be fetched: the bridge would be "+ + "reading its own internal network on a stranger's instructions, and whatever "+ + "answered would decide whether that stranger's account is erased") + assert.False(t, deleted, "and above all it must not produce the deleted verdict") +} + +// asSplitTransport serves the DID document locally and sends everything else — +// the PDS leg, which is the one under test — through the production guard. +type asSplitTransport struct { + directoryHost string + directory http.RoundTripper + guarded http.RoundTripper +} + +func (rt *asSplitTransport) RoundTrip(req *http.Request) (*http.Response, error) { + if !strings.HasPrefix(req.URL.Path, asStatusPath) { + return rt.directory.RoundTrip(req) + } + return rt.guarded.RoundTrip(req) +} diff --git a/internal/consume/federation.go b/internal/consume/federation.go index 4c30da0..d2ab9a8 100644 --- a/internal/consume/federation.go +++ b/internal/consume/federation.go @@ -77,13 +77,19 @@ func (d *Dispatcher) handleFederation(ctx context.Context, tx *sql.Tx, did strin } // STOP MEANS TWO FACTS AT ONCE: nothing new goes out, and nothing already - // queued goes out either. The cancellation answers the second, and it runs - // FIRST — before the destructive tier below — because that tier ENQUEUES - // the withdrawal, and a cancellation of "this actor's pending work" that ran - // afterwards would cancel the Delete{Person} it just queued. The order is - // not a preference: the two statements are the same predicate pointed at - // different moments. - cancelled, err := d.deliveries.CancelForActorTx(ctx, tx, did) + // queued goes out either. The cancellation answers the second — and it + // cancels only what PUBLISHES, never a retraction. + // + // That exemption is what makes this replay-safe. The destructive tier below + // enqueues its withdrawal (a Delete and some Undos) on its own transaction, + // which can commit while this one later rolls back; the record then replays + // and reaches this line again. A sweeping cancel here would cancel the + // erasure the previous attempt committed, and nothing would repair it — the + // delivery insert returns the standing row rather than reviving it, and the + // votes are already retracted, so the re-run finds nothing to enumerate. The + // user would be tombstoned, the log would say the purge applied, and no peer + // would ever have been told. + cancelled, err := d.deliveries.CancelOutwardForActorTx(ctx, tx, did) if err != nil { return fmt.Errorf("cancel queued deliveries for %s: %w", did, err) } @@ -134,9 +140,9 @@ func (d *Dispatcher) deleteRemoteContent(ctx context.Context, did string) error slog.String("did", did)) return nil } - // Reached exactly once per applied record: the rev gate rejects a replay - // before this handler runs, and asking peers a second time to delete - // content the user may have since re-enabled is unrecoverable. + // NOT once per record, despite the gate: this call can run again on the + // replay its own doc describes, so the purge has to be idempotent rather + // than merely rare (it is — see outbound.Purger). if err := d.remoteDeleter.DeleteRemoteContent(ctx, did); err != nil { return fmt.Errorf("delete remote content for %s: %w", did, err) } @@ -151,11 +157,36 @@ func (d *Dispatcher) deleteRemoteContent(ctx context.Context, did string) error // ability to federate from now on; the work they cancelled by asking us to stop // was withdrawn at their request, and re-sending it would publish on their // behalf something they had already taken back. +// A WITHDRAWN IDENTITY IS NOT RESTORED BY EITHER. The destructive tier asked +// every peer to delete this user's content and its actor document answers 410 +// forever; re-enabling afterwards would sign new posts, comments and votes as +// somebody peers were explicitly told is gone. Both halves of the restore are +// refused at their own store — SetEnabled cannot re-enable a tombstoned actor, +// 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. func (d *Dispatcher) restoreDefaultFederation(ctx context.Context, tx *sql.Tx, did string) error { if err := d.prefs.Delete(ctx, did); err != nil { return fmt.Errorf("clear federation preference for %s: %w", did, err) } - return d.mirrorActorEnabled(ctx, tx, did, true) + if err := d.mirrorActorEnabled(ctx, tx, did, true); err != nil { + return err + } + + // 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 + case err != nil: + return fmt.Errorf("read federation preference for %s: %w", did, err) + case pref.PurgedAt != nil: + d.logger.Warn("re-enable refused: this identity was withdrawn and the withdrawal is irreversible", + slog.String("did", did), slog.Time("purged_at", *pref.PurgedAt)) + } + return nil } // mirrorActorEnabled updates the ap_actors mirror IF the actor exists. A diff --git a/internal/consume/rev_gate.go b/internal/consume/rev_gate.go index ba775d6..4f86311 100644 --- a/internal/consume/rev_gate.go +++ b/internal/consume/rev_gate.go @@ -175,10 +175,20 @@ func (g *RevGate) Advance(ctx context.Context, uri, rev string) error { // skips — so no check→write window exists for a stale copy to sneak through, // no matter how long apply takes. // -// Deadlock note: apply's writes go through repository methods on their own -// connections, which is deliberate and safe — the gate transaction touches -// ONLY jetstream_record_revs, a table no repository write path ever touches, -// so the gate row lock acts as a pure per-record mutex around apply. +// DEADLOCK NOTE, and it is no longer as simple as it once was. The gate +// transaction originally touched ONLY jetstream_record_revs, which made the +// 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 rule that replaces the old proof: A HANDLER WRITING ON THIS TRANSACTION +// MUST NOT THEN CALL SOMETHING THAT OPENS A SECOND TRANSACTION TOUCHING THE +// SAME ROWS. It cannot release what it holds without committing, and the +// 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. // // An apply error — or a panic, which the deferred rollback covers equally — // releases the claim WITHOUT advancing, so the connector's retry/redrive diff --git a/internal/db/migrations/030_federation_pref_purged.sql b/internal/db/migrations/030_federation_pref_purged.sql new file mode 100644 index 0000000..4645299 --- /dev/null +++ b/internal/db/migrations/030_federation_pref_purged.sql @@ -0,0 +1,32 @@ +-- +goose Up +-- Task 17d review (P1-b): telling a purge that was REQUESTED from one that +-- actually HAPPENED. +-- +-- The terminal tier records the preference BEFORE it asks peers to delete +-- anything, deliberately: a Delete no peer can un-honour must not be sent from +-- an intent that a crash could lose. But that leaves a window where the row says +-- "this user is withdrawn" and nothing was withdrawn — and if the purge FAILS +-- while the user REACTIVATES, the next confirmation returns live, the event is +-- handled, and a live account is left permanently opted out by a decision that +-- was never carried out. Nothing clears it, because nothing knows the difference. +-- +-- purged_at is that difference, and it is on THIS row rather than inferred from +-- the actor's tombstone because this is the row a later re-enable would clear: +-- the check and the thing being checked belong together, and an actor may not +-- exist at all (a DID that never federated anything has nothing to tombstone). +-- +-- NULL — requested. Nothing irreversible happened, so a confirmed-live +-- verdict may clear it and the user carries on. +-- set — the withdrawal COMMITTED. Peers were asked to delete, and no +-- later record, re-enable or reactivation may erase the evidence +-- or resume federating: there is nothing to come back to. +-- +-- Nothing clears this column. It is the same terminality as ap_actors +-- .tombstoned_at, recorded where the opposite decision would be written. +ALTER TABLE federation_prefs ADD COLUMN purged_at TIMESTAMPTZ; + +-- +goose Down +-- The column goes; the rows STAY. Dropping the preference rows to satisfy a +-- narrower schema would resume federating for accounts whose content peers were +-- already told to delete — the one outcome this whole tier exists to prevent. +ALTER TABLE federation_prefs DROP COLUMN IF EXISTS purged_at; diff --git a/internal/ingest/optout_retraction_test.go b/internal/ingest/optout_retraction_test.go new file mode 100644 index 0000000..c78825d --- /dev/null +++ b/internal/ingest/optout_retraction_test.go @@ -0,0 +1,181 @@ +package ingest + +import ( + "context" + "database/sql" + "encoding/json" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/consume" +) + +// TASK 17d REVIEW — WHAT A STOPPED USER IS STILL OWED. +// +// "Stop federating for me" and "forget what you already sent" are different +// instructions, and only the first one is ever asked for. The worker has said so +// since task 15 — a Delete or an Undo skips the consent recheck, because a +// retraction is the ONLY way an opted-out user takes down what is already +// published — but the queue-side cancellation is one layer up, where nothing +// observed it. +// +// Two consequences, and the second is the one that eats the whole tier: +// +// 1. A user who deletes a post and THEN opts out has their pending self-delete +// cancelled, so the post stays on Lemmy forever. They asked to be forgotten +// and got the opposite of both halves. +// 2. THE DESTRUCTIVE TIER CANCELS ITSELF ON REPLAY. The purge commits on its own +// transaction; the rev gate commits later. A failure in between replays the +// record — and the replay's cancellation runs FIRST, sweeping the +// Delete{Person} fan-out and the vote Undos the previous attempt just +// committed. Nothing repairs it: the delivery insert returns the standing +// (cancelled) row by design, and the votes are already flipped to 'undone' so +// the re-run enumerates nothing to retract. The actor is tombstoned, the +// preference is recorded, the log says "destructive opt-out applied", and not +// one peer was ever told. +// +// (2) is invisible from inside a single pass: every row exists, every counter +// moves. Only replaying the record — which is what a rolled-back gate DOES — +// shows it. + +const ( + orSelfDeleteRKey = "3lzorpost00001" + orOptOutRev = "3lzorrev000001" + orPurgeRev = "3lzorrev000002" + orPurgeRKey = "3lzorpost00002" +) + +// TestASelfDeleteStillGoesOutAfterTheAuthorOptsOut is (1). +// +// The delete is already in the queue when the opt-out lands. Cancelling it is +// the one cancellation that leaves MORE of the user's content on the fediverse +// than doing nothing would have. +func TestASelfDeleteStillGoesOutAfterTheAuthorOptsOut(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + // --- GIVEN: an accepted post the author then takes down. + admitPost(t, world, mtAuthorDID, orSelfDeleteRKey, world.communityADID, + "3lzorrev000010", 1_775_000_050_000_001) + require.NoError(t, world.dispatcher.HandleEvent(ctx, + postDeleteEvent(t, mtAuthorDID, orSelfDeleteRKey, "3lzorrev000011", 1_775_000_050_000_002))) + + deleteActivity := activityOfKind(t, h.db, mtAuthorDID, "Delete") + require.NotEmpty(t, deleteActivity, + "precondition: the author's own delete really did enqueue a retraction") + require.Equal(t, []string{"pending"}, deliveryStatesForActivity(t, h.db, deleteActivity), + "precondition: and it has not gone out yet — this is the window the opt-out lands in") + + // --- WHEN: they opt out. Soft tier. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, mtAuthorDID, orOptOutRev, "create", false, false, 1_775_000_051_000_001))) + + // --- THEN: the retraction survives. + assert.Equal(t, []string{"pending"}, deliveryStatesForActivity(t, h.db, deleteActivity), + "a Delete is how an opted-out user takes down what is ALREADY federated. Cancelling "+ + "it leaves the post standing on Lemmy forever — the user asked to stop federating "+ + "and the queue answered by keeping their content published. The worker has always "+ + "exempted retractions from the consent recheck; a cancellation one layer up "+ + "reverses that decision where the worker can no longer see it") +} + +// TestAReplayedDestructiveOptOutDoesNotCancelItsOwnWithdrawal is (2). +// +// The gate rollback is spelled the way it actually happens: the purge committed, +// the gate did not, so no rev row survives and the SAME record is delivered +// again. Everything the second pass does must leave the first pass's withdrawal +// deliverable. +func TestAReplayedDestructiveOptOutDoesNotCancelItsOwnWithdrawal(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + // --- GIVEN: content on two instances and one vote a peer still holds. + admitPost(t, world, mtAuthorDID, orPurgeRKey, world.communityADID, + "3lzorrev000020", 1_775_000_052_000_001) + h.subscribeCommunityURL(odFarCommunity, odFarName) + admitPost(t, world, mtAuthorDID, "3lzorpost00003", testDIDFor(odFarName, "lemmy.zip"), + "3lzorrev000021", 1_775_000_052_000_002) + liveVote := seedDeliveredVote(t, h.db, mtAuthorDID, mtPostATURI, world.communityADID, "delivered") + + // --- WHEN: the destructive record is applied... + event := odFederationEvent(t, mtAuthorDID, orPurgeRev, "create", false, true, 1_775_000_053_000_001) + require.NoError(t, world.dispatcher.HandleEvent(ctx, event)) + + deleteActivity := activityOfKind(t, h.db, mtAuthorDID, "Delete") + require.NotEmpty(t, deleteActivity, "precondition: the withdrawal was enqueued") + inboxes := deliveryStatesForActivity(t, h.db, deleteActivity) + require.Len(t, inboxes, 2, "precondition: fanned out to both instances") + require.Equal(t, 1, undoActivitiesFor(t, h.db, mtAuthorDID), + "precondition: and the live vote's Undo with it") + require.Equal(t, "undone", voteState(t, h.db, liveVote), + "precondition: the vote is already flipped, so a re-run enumerates nothing — which is "+ + "exactly why a cancelled Undo is unrecoverable") + + // --- ...and the gate transaction rolls back. Nothing marks the record + // applied, so the connector redelivers it and the handler runs again. + clearRevGate(t, h.db, mtAuthorDID) + require.NoError(t, world.dispatcher.HandleEvent(ctx, event), + "the replay is the ordinary consequence of a rolled-back gate, not an error") + + // --- THEN: the withdrawal is still deliverable. + states := deliveryStatesForActivity(t, h.db, deleteActivity) + require.Len(t, states, 2, + "the replay must not lose an inbox: the standing row wins, so the fan-out is neither "+ + "duplicated nor dropped") + for i, state := range states { + assert.Equal(t, "pending", state, + "delivery %d of %d: a replayed opt-out must not cancel the Delete{Person} the "+ + "previous pass committed. The cancellation runs before the purge and cannot "+ + "tell its own withdrawal from the user's ordinary queued work — and nothing "+ + "repairs it, because the re-enqueue returns the standing cancelled row and the "+ + "votes are already retracted. The result is an actor tombstoned, a preference "+ + "recorded, a log line saying it was applied, and zero peers told", i+1, len(states)) + } + + undoActivity := activityOfKind(t, h.db, mtAuthorDID, "Undo") + require.NotEmpty(t, undoActivity) + for _, state := range deliveryStatesForActivity(t, h.db, undoActivity) { + assert.Equal(t, "pending", state, + "and the vote Undo with it: the peer holds a vote from an actor we have just "+ + "tombstoned, and this delivery is the only thing that will ever tell them") + } +} + +// postDeleteEvent is the author's own take-down: a delete commit, which carries +// no record and no CID — everything the Delete{Page} needs is read from the +// stored outbound state. +func postDeleteEvent(t *testing.T, did, rkey, rev string, timeUS int64) *consume.JetstreamEvent { + t.Helper() + frame := fmt.Sprintf( + `{"did":%q,"time_us":%d,"kind":"commit","commit":{"rev":%q,"operation":"delete",`+ + `"collection":"social.coves.community.postv2","rkey":%q}}`, + did, timeUS, rev, rkey) + var event consume.JetstreamEvent + require.NoError(t, json.Unmarshal([]byte(frame), &event), "the frame must be valid wire JSON") + return &event +} + +// clearRevGate deletes the DID's rev-gate rows, which is the state a ROLLED-BACK +// gate transaction leaves behind: the handler's own writes may have committed +// separately (the purge takes its own transaction), and nothing records that the +// record was applied, so the connector redelivers it. +func clearRevGate(t *testing.T, db *sql.DB, did string) { + t.Helper() + _, err := db.ExecContext(context.Background(), + `DELETE FROM jetstream_record_revs WHERE record_uri LIKE 'at://' || $1 || '/%'`, did) + require.NoError(t, err) +} + +// deliveryStatesForActivity lists every delivery state for one activity, in a +// stable order. +func deliveryStatesForActivity(t *testing.T, db *sql.DB, activityID string) []string { + t.Helper() + return queryStrings(t, db, + `SELECT state FROM outbound_deliveries WHERE activity_id = $1 ORDER BY target_inbox`, + activityID) +} diff --git a/internal/optout/terminator.go b/internal/optout/terminator.go index 66f5d22..8fc2526 100644 --- a/internal/optout/terminator.go +++ b/internal/optout/terminator.go @@ -120,7 +120,7 @@ func (t *Terminator) TerminateAccount(ctx context.Context, did string) error { // a confirm that is broken indistinguishable from one that is working. t.logger.Info("account deletion not confirmed; taking no destructive action", slog.String("did", did)) - return nil + return t.clearStaleRequest(ctx, did) } // RECORDED BEFORE ANYTHING IS SENT. Peers that honour a Delete cannot @@ -149,7 +149,40 @@ func (t *Terminator) TerminateAccount(ctx context.Context, did string) error { if err := t.deleter.DeleteRemoteContent(ctx, did); err != nil { return fmt.Errorf("delete remote content for %s: %w", did, err) } + // STAMPED ONLY NOW. Everything above this line is a REQUEST — recorded first + // so a crash could not lose the user's intent — and a request is something a + // live account may still withdraw. This marks the moment it stopped being + // one: peers have been asked to delete, and nothing after this may clear the + // preference or bring the identity back. + if err := t.prefs.MarkPurged(ctx, did); err != nil { + return fmt.Errorf("mark account purge committed for %s: %w", did, err) + } t.logger.Info("account confirmed deleted; remote content withdrawal requested", slog.String("did", did)) return nil } + +// clearStaleRequest withdraws a preference THIS TIER wrote for a deletion that +// never happened. +// +// The window it closes: an account is reported deleted, the preference is +// recorded, the purge FAILS, and before the retry the user reactivates. The +// confirmation now returns live, so the destructive path never runs again — and +// without this the account-sourced "disabled" row would stand forever, blocking +// a live user with no purge having committed and nothing left to clear it. +// +// It can only ever remove a request. The store refuses to clear a user's own +// opt-out (theirs to keep) or a purge that committed (peers were already told), +// so the narrow case is narrow by construction rather than by this caller +// getting the predicate right. +func (t *Terminator) clearStaleRequest(ctx context.Context, did string) error { + cleared, err := t.prefs.ClearRequestedPurge(ctx, did) + if err != nil { + return fmt.Errorf("clear stale deletion request for %s: %w", did, err) + } + if cleared { + t.logger.Info("account is live again; the recorded deletion request was withdrawn", + slog.String("did", did)) + } + return nil +} diff --git a/internal/outbound/enqueuer.go b/internal/outbound/enqueuer.go index c70d8b5..103a9d4 100644 --- a/internal/outbound/enqueuer.go +++ b/internal/outbound/enqueuer.go @@ -188,7 +188,7 @@ func (e *Enqueuer) EnqueueActivity(ctx context.Context, tx *sql.Tx, actorDID, or // Zero targets is NOT an error. An actor whose content never reached anyone has // nothing to withdraw, and failing here would turn "nothing to do" into an event // that retries forever. -func (e *Enqueuer) EnqueueFanOut(ctx context.Context, tx *sql.Tx, actorDID string, +func (e *Enqueuer) EnqueueFanOut(ctx context.Context, tx *sql.Tx, actorDID, parentATURI string, intent consume.Intent, targets []store.DeliveryTarget) error { if tx == nil { @@ -206,10 +206,11 @@ func (e *Enqueuer) EnqueueFanOut(ctx context.Context, tx *sql.Tx, actorDID strin return fmt.Errorf("translate intent %s: %w", intent.ActivityID(), err) } if _, err := e.activities.InsertTx(ctx, tx, store.OutboundActivity{ - ActivityID: intent.ActivityID(), - ActorDID: actorDID, - Kind: translated.Kind, - Payload: translated.Payload, + ActivityID: intent.ActivityID(), + ActorDID: actorDID, + Kind: translated.Kind, + Payload: translated.Payload, + ParentATURI: parentATURI, }); err != nil { return fmt.Errorf("insert outbound activity %s: %w", intent.ActivityID(), err) } @@ -226,6 +227,18 @@ func (e *Enqueuer) EnqueueFanOut(ctx context.Context, tx *sql.Tx, actorDID strin return nil } +// Inbox resolves a community's shared inbox through the same cached resolver +// EnqueueActivity uses. +// +// It is exported for ONE caller with one reason: the destructive tier must +// resolve every target BEFORE it opens its transaction. Resolving inside would +// hold locks on outbound_deliveries for as long as N remote actor fetches take, +// and everything else that touches those rows — the worker, another opt-out, a +// test harness truncating between cases — waits behind a network call. +func (e *Enqueuer) Inbox(ctx context.Context, communityAPID string) (string, error) { + return e.inboxes.ResolveInbox(ctx, communityAPID) +} + // objectMapping derives the bridge-origin ap_objects mapping for an intent that // produces a servable object (a comment or post create/update). Votes have no // object; a self-delete's object was mapped on its create. ok=false means no diff --git a/internal/outbound/fanout_test.go b/internal/outbound/fanout_test.go index 698bffe..fd69e02 100644 --- a/internal/outbound/fanout_test.go +++ b/internal/outbound/fanout_test.go @@ -176,7 +176,9 @@ func TestEnqueuer_ARedeliveredActivityToTheSameInboxStaysOneRow(t *testing.T) { "second row would send the peer a duplicate it has to recognise and discard") } -// consumeIntentCompileGuard keeps the fan-out fixture honest about the seam it -// stands in for: whatever intent the destructive tier introduces, it reaches -// these same two tables through this same method. -var _ consume.Intent = consume.CommentIntent{} +// The intent the destructive tier actually introduces. CommentIntent (which +// every test above already exercises) proves nothing about it: PersonDeleteIntent +// is the one activity in the system addressed to MANY inboxes, and it is the +// reason EnqueueFanOut exists at all. Naming it here is what makes the fixture's +// claim to stand in for the destructive tier true. +var _ consume.Intent = consume.PersonDeleteIntent{} diff --git a/internal/outbound/purge.go b/internal/outbound/purge.go index 2c55dbe..4ffeabe 100644 --- a/internal/outbound/purge.go +++ b/internal/outbound/purge.go @@ -19,10 +19,27 @@ import ( // // 1. Delete{Person, removeData:true} to every inbox the actor's content // reached. The instances it misses keep serving those posts forever. +// // 2. An Undo for every vote a peer still holds. Tidepool's aggregate is the // FEDIVERSE-ONLY tally and the reseed SUBTRACTS live delivered votes from // the origin's API count (task 17b), so a standing vote from a withdrawn // actor is a number the reseed keeps subtracting from a score readers see. +// +// "Standing" is wider than delivered_state='delivered', deliberately: a vote +// whose delivery is HELD FOR SETTLEMENT is one the peer ALREADY ACCEPTED +// with only our bookkeeping outstanding, and its settlement lands after this +// withdrawal — so enumerating the settled rows alone would leave a real vote +// on a real instance with nothing left to notice (ListStandingForActor). +// The retraction is what makes that safe in both directions: `undone` is +// terminal in SetDeliveredState, so the late settlement cannot put the vote +// back. +// +// RESIDUAL, and it is the one nothing here can close: a delivery CLAIMED and +// mid-POST at this moment is indistinguishable from one that will never be +// sent. If that POST lands after the purge, the peer holds a vote we never +// retracted. Detecting it needs the peer's own state, which is +// reconciliation — 17e's. +// // 3. The actor document stops resolving — 410 Gone. // // IRREVERSIBLE, and reached only from an explicit enabled=false + @@ -42,6 +59,7 @@ type Purger struct { enqueuer *Enqueuer deliveries store.OutboundDeliveries votes store.OutboundVotes + prefs store.FederationPrefs actors store.APActors communities store.Communities logger *slog.Logger @@ -56,6 +74,7 @@ func NewPurger(db *sql.DB, userOrigin string, enqueuer *Enqueuer) *Purger { enqueuer: enqueuer, deliveries: store.NewOutboundDeliveries(db), votes: store.NewOutboundVotes(db), + prefs: store.NewFederationPrefs(db), actors: store.NewAPActors(db), communities: store.NewCommunities(db), logger: slog.Default(), @@ -93,13 +112,20 @@ func (p *Purger) DeleteRemoteContent(ctx context.Context, did string) error { return errors.NewValidationError("did", "must not be empty") } - // Both reads happen BEFORE the transaction opens: they are the inputs, and - // holding a transaction open across them buys nothing. + // EVERY READ AND EVERY REMOTE LOOKUP HAPPENS BEFORE THE TRANSACTION OPENS. + // Resolving a vote's community can touch the network (the inbox resolver + // caches, but a cold entry fetches an actor document), and holding a + // transaction open across N remote fetches keeps locks for as long as the + // slowest peer takes to answer. targets, err := p.deliveries.DistinctInboxesForActor(ctx, did) if err != nil { return err } - liveVotes, err := p.votes.ListDeliveredForActor(ctx, did) + liveVotes, err := p.votes.ListStandingForActor(ctx, did) + if err != nil { + return err + } + retractable, err := p.addressableVotes(ctx, did, liveVotes) if err != nil { return err } @@ -110,7 +136,7 @@ func (p *Purger) DeleteRemoteContent(ctx context.Context, did string) error { } defer func() { _ = tx.Rollback() }() - if err := p.enqueuer.EnqueueFanOut(ctx, tx, did, consume.PersonDeleteIntent{ + if err := p.enqueuer.EnqueueFanOut(ctx, tx, did, "", consume.PersonDeleteIntent{ ActorDID: did, // seq 0: a withdrawal happens once per identity, and the id must be the // SAME string on every retry so a peer recognises the redelivery as the @@ -120,10 +146,24 @@ func (p *Purger) DeleteRemoteContent(ctx context.Context, did string) error { return err } - if err := p.undoLiveVotes(ctx, tx, did, liveVotes); err != nil { + if err := p.undoLiveVotes(ctx, tx, did, retractable); err != nil { return err } + // The preference this withdrawal was asked for stops being a REQUEST here, + // inside the same transaction as the withdrawal itself. Both doors reach + // this code — an explicit deleteRemote record and a confirmed account + // deletion — and only the marked row is protected from being cleared by a + // later re-enable, so marking it anywhere but here would leave one door + // open. A missing preference does not fail the purge: it is a hole in the + // caller's ordering, not a reason to abandon a withdrawal already enqueued. + if err := p.prefs.MarkPurgedTx(ctx, tx, did); err != nil && !errors.IsNotFound(err) { + return fmt.Errorf("mark purge committed for %s: %w", did, err) + } else if errors.IsNotFound(err) { + p.logger.Warn("purge committed with no federation preference to mark", + slog.String("did", did)) + } + // LAST, because it is the step that stops the actor being servable and the // enqueues above resolve that actor. It is also the one an operator will // read as "this user is gone", so it must not be true before the withdrawal @@ -138,7 +178,8 @@ func (p *Purger) DeleteRemoteContent(ctx context.Context, did string) error { p.logger.Info("destructive opt-out applied; the withdrawal is irreversible", slog.String("did", did), slog.Int("inboxes", len(targets)), - slog.Int("votes_retracted", len(liveVotes))) + slog.Int("votes_retracted", len(retractable)), + slog.Int("votes_unaddressable", len(liveVotes)-len(retractable))) return nil } @@ -153,17 +194,9 @@ func (p *Purger) DeleteRemoteContent(ctx context.Context, did string) error { // The cost is stated plainly — if the Undo never lands, we have stopped counting // a vote the peer may still hold — and it is the right side to err on, because // the alternative subtracts forever on behalf of somebody who is gone. -func (p *Purger) undoLiveVotes(ctx context.Context, tx *sql.Tx, did string, votes []store.OutboundVote) error { - for i := range votes { - vote := votes[i] - // The community resolves BEFORE the row is bumped, so a failed lookup - // rolls back rather than leaving a bumped seq behind an Undo that was - // never enqueued (consume's vote delete draws the same line). - community, err := p.communities.GetByDID(ctx, vote.CommunityDID) - if err != nil { - return fmt.Errorf("resolve community %s for vote %s: %w", - vote.CommunityDID, vote.VoteATURI, err) - } +func (p *Purger) undoLiveVotes(ctx context.Context, tx *sql.Tx, did string, votes []retractableVote) error { + for _, retractable := range votes { + vote := retractable.vote // One statement bumps the seq — the Undo is the next activity and its id // must not collide with the Like's — and flips the state. CurrentActivityID @@ -175,7 +208,7 @@ func (p *Purger) undoLiveVotes(ctx context.Context, tx *sql.Tx, did string, vote return fmt.Errorf("retract vote %s: %w", vote.VoteATURI, err) } - if err := p.enqueuer.EnqueueActivity(ctx, tx, did, did, vote.SubjectATURI, consume.VoteIntent{ + if err := p.enqueuer.EnqueueFanOut(ctx, tx, did, vote.SubjectATURI, consume.VoteIntent{ Op: consume.OperationUndo, VoteATURI: vote.VoteATURI, SubjectAPID: vote.SubjectAPID, @@ -184,10 +217,66 @@ func (p *Purger) undoLiveVotes(ctx context.Context, tx *sql.Tx, did string, vote Direction: vote.Direction, ID: consume.ActivityID(p.userOrigin, vote.VoteATURI, consume.OperationUndo, bumped.ActivitySeq), InnerActivityID: vote.CurrentActivityID, - CommunityAPID: community.APGroupID, - }); err != nil { + CommunityAPID: retractable.communityAPID, + }, []store.DeliveryTarget{{ + Inbox: retractable.inbox, + OrderingKey: retractable.communityAPID, + }}); err != nil { return err } } return nil } + +// retractableVote is a live vote paired with the community AP id its Undo is +// addressed to — resolved before the transaction opens, so no remote lookup +// happens under a lock. +type retractableVote struct { + vote store.OutboundVote + communityAPID string + inbox string +} + +// addressableVotes resolves the addressing for each live vote and DROPS the ones +// that cannot be addressed at all. +// +// A community that has been deleted or unfollowed since the vote was cast has no +// row, and there is no inbox to send its Undo to. Treating that as fatal would +// hold the ENTIRE erasure hostage to one vote: the transaction rolls back, the +// Delete{Person} fan-out with it, the event replays into the same missing row +// and eventually dead-letters — a user's whole withdrawal lost to a community +// that no longer exists. +// +// A genuine storage error is still fatal. "This community is gone" and "the +// database did not answer" are different facts, and only the first one is an +// answer. +func (p *Purger) addressableVotes(ctx context.Context, did string, votes []store.OutboundVote) ([]retractableVote, error) { + out := make([]retractableVote, 0, len(votes)) + for i := range votes { + vote := votes[i] + community, err := p.communities.GetByDID(ctx, vote.CommunityDID) + if errors.IsNotFound(err) { + p.logger.Warn("purge: a live vote's community is gone; its Undo cannot be addressed", + slog.String("did", did), slog.String("vote", vote.VoteATURI), + slog.String("community_did", vote.CommunityDID)) + continue + } + if err != nil { + return nil, fmt.Errorf("resolve community %s for vote %s: %w", + vote.CommunityDID, vote.VoteATURI, err) + } + // The INBOX is resolved here too, and that is the whole reason this + // function exists outside the transaction: resolving one can fetch a + // remote actor document, and doing that with locks held on + // outbound_deliveries stalls the worker — and anything else touching + // those rows — for as long as the slowest peer takes to answer. + inbox, err := p.enqueuer.Inbox(ctx, community.APGroupID) + if err != nil { + return nil, fmt.Errorf("resolve inbox for %s: %w", community.APGroupID, err) + } + out = append(out, retractableVote{ + vote: vote, communityAPID: community.APGroupID, inbox: inbox, + }) + } + return out, nil +} diff --git a/internal/outbound/purge_held_vote_test.go b/internal/outbound/purge_held_vote_test.go new file mode 100644 index 0000000..3daccea --- /dev/null +++ b/internal/outbound/purge_held_vote_test.go @@ -0,0 +1,224 @@ +package outbound + +import ( + "context" + "database/sql" + stderrors "errors" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" + "tidepool/internal/testutil" +) + +// TASK 17d REVIEW — THE VOTE THE ERASURE NEVER SEES. +// +// The purge enumerates what to retract from outbound_votes where delivered_state +// = 'delivered', read BEFORE its transaction opens. A vote whose delivery is +// HELD FOR SETTLEMENT is invisible to that read: the peer HAS the Like — the +// POST was confirmed — and the ledger still says 'pending' because only the +// local settlement failed. +// +// The hold is not a lost cause; it is the opposite. The worker comes back, +// resumes at the settlement, and flips the row to 'delivered'. So the sequence +// that this test drives is the ordinary one: +// +// POST confirmed → ledger write fails → row HELD, vote 'pending' +// purge runs → enumerates nothing → no Undo owed +// worker resumes → settlement lands → vote 'delivered', actor tombstoned +// +// and it ends with a vote standing on a peer, attributed to an actor this bridge +// has told the world is gone, with no Undo ever enqueued and nothing left that +// will ever notice — the purge is terminal and never re-runs. 17b's reseed keeps +// subtracting it from the served score forever, and the erasure has silently +// failed at the one thing it exists to do. +// +// THE ORDERING GAP IS THE BUG, not the hold. Whatever closes it — enumerating +// held deliveries too, retracting inside the transaction, or a settlement that +// checks for a tombstoned actor — the property below is the same, and it is the +// property a user asked for. + +const pvVoteATURI = "at://" + wActorDID + "/social.coves.feed.vote/3lzpvvote001" + +// purgeTestDB is workerTestDB plus the community table the purge resolves each +// vote's target through. +func purgeTestDB(t *testing.T) *sql.DB { + t.Helper() + conn := workerTestDB(t) + testutil.Truncate(t, conn, "communities") + _, err := store.NewCommunities(conn).UpsertCommunity(context.Background(), store.Community{ + APGroupID: wCommunityAPID, + DID: wCommunityDID, + PreferredUsername: "tech", + Instance: "lemmy.world", + FollowState: store.FollowStateAccepted, + }) + require.NoError(t, err) + return conn +} + +// pvFaultyVotes fails the delivery-success ledger write on demand, which is what +// puts a confirmed POST into the held state. +type pvFaultyVotes struct { + store.OutboundVotes + failSet bool +} + +func (f *pvFaultyVotes) SetDeliveredState(ctx context.Context, voteATURI string, state store.DeliveredState) error { + if f.failSet { + return stderrors.New("ledger write failed") + } + return f.OutboundVotes.SetDeliveredState(ctx, voteATURI, state) +} + +// heldVoteWorld is one persona's vote, confirmed on the wire, with its ledger +// write faulted — the state both tests below start from. +type heldVoteWorld struct { + conn *sql.DB + worker *Worker + votes store.OutboundVotes + faulty *pvFaultyVotes + purger *Purger + likeID string +} + +func newHeldVoteWorld(t *testing.T) *heldVoteWorld { + t.Helper() + conn := purgeTestDB(t) + ctx := context.Background() + seedWorkerActor(t, conn, true, false) + + likeID := "https://coves.social/ap/activity/" + repeatHex64("Like") + votes := store.NewOutboundVotes(conn) + _, err := votes.Upsert(ctx, store.OutboundVote{ + VoteATURI: pvVoteATURI, + ActorDID: wActorDID, + SubjectATURI: "at://" + wCommunityDID + "/social.coves.community.postv2/3lzpost", + SubjectAPID: "https://lemmy.world/post/1", + CommunityDID: wCommunityDID, + Direction: "up", + CurrentActivityID: likeID, + }) + require.NoError(t, err) + seedDelivery(t, conn, "Like", "", []byte(fmt.Sprintf( + `{"id":%q,"type":"Like","actor":%q,"object":"https://lemmy.world/post/1"}`, likeID, wActorID))) + + // The Like reaches the peer and the ledger write does not commit. + faulty := &pvFaultyVotes{OutboundVotes: votes, failSet: true} + sender := &fakeSender{} + w := newWorker(t, conn, sender, func(o *WorkerOptions) { o.Votes = faulty }) + _, err = w.DeliverNext(ctx) + require.NoError(t, err) + require.Equal(t, 1, sender.count(), + "precondition: the vote really is AT the peer — everything here is about what happens "+ + "after the wire said yes") + + held := getDelivery(t, conn, likeID) + require.Equal(t, store.DeliveryStatePending, held.State) + require.Equal(t, store.DeliveryHeldForSettlement, held.LastErrorClass, + "precondition: the delivery is held for settlement, so the worker WILL come back to it") + stored, err := votes.GetByATURI(ctx, pvVoteATURI) + require.NoError(t, err) + require.Equal(t, store.DeliveredStatePending, stored.DeliveredState, + "precondition: and the ledger row still says pending, which is the whole gap") + + enqueuer, err := NewEnqueuer(EnqueuerOptions{ + DB: conn, + Translator: NewTranslator("https://coves.social"), + Inboxes: staticInbox{inbox: wInbox}, + Actors: store.NewAPActors(conn), + UserOrigin: "https://coves.social", + }) + require.NoError(t, err) + + return &heldVoteWorld{ + conn: conn, worker: w, votes: votes, faulty: faulty, + purger: NewPurger(conn, "https://coves.social", enqueuer), likeID: likeID, + } +} + +// settle lets the held delivery finish the job it was held for, and asserts it +// really did — on the DELIVERY, which is the thing the hold is about. Asserting +// on the vote row instead would make the precondition depend on the very value +// the terminality test exists to pin. +func (w *heldVoteWorld) settle(t *testing.T) { + t.Helper() + w.faulty.failSet = false + _, err := w.worker.DeliverNext(context.Background()) + require.NoError(t, err) + require.Equal(t, store.DeliveryStateDelivered, getDelivery(t, w.conn, w.likeID).State, + "precondition: the held settlement completed and the delivery reached its terminal "+ + "state — this is the recovery working, not a fault") +} + +func TestPurge_DoesNotLeaveAHeldVoteStandingOnAPeer(t *testing.T) { + world := newHeldVoteWorld(t) + ctx := context.Background() + conn := world.conn + + // --- WHEN: the user's account is withdrawn, and then the settlement the + // hold was waiting for lands. + require.NoError(t, world.purger.DeleteRemoteContent(ctx, wActorDID)) + world.settle(t) + + // --- THEN: the peer must not be left holding a vote from a withdrawn actor. + assert.Equal(t, 1, activitiesOfKind(t, conn, "Undo"), + "the erasure owes this peer an Undo. The purge read outbound_votes BEFORE its "+ + "transaction and the row said 'pending', so it enumerated nothing — but the peer "+ + "had already accepted the Like, and the settlement that followed flipped the row to "+ + "'delivered' for an actor we have since tombstoned. Nothing re-runs a purge: the "+ + "vote stands on that instance forever, and 17b's reseed keeps subtracting it from a "+ + "score readers see. A hold is a delivery that SUCCEEDED, so the enumeration that "+ + "decides what a withdrawal owes cannot read only the rows whose bookkeeping "+ + "happened to finish first") +} + +// TestPurge_ARetractedVoteIsNotResurrectedByALateSettlement is the OTHER half, +// and it is the half the enumeration fix cannot reach. +// +// Enumerating the held vote is what gets the Undo sent. It does nothing about +// what happens NEXT: the worker still comes back, the held settlement still +// runs, and SetDeliveredState still writes 'delivered' over the 'undone' the +// purge just recorded. The Undo is on the wire and the ledger says the vote is +// live — so 17b's reseed subtracts it from a served score forever, and an +// operator auditing a withdrawn identity sees votes it supposedly still holds. +// +// `undone` is the one value in this column that records OUR decision rather +// than the message's progress. A decision cannot be overwritten by a later fact +// about the thing it was a decision ABOUT, and this is the one place the two +// arrive out of order. +func TestPurge_ARetractedVoteIsNotResurrectedByALateSettlement(t *testing.T) { + world := newHeldVoteWorld(t) + ctx := context.Background() + + require.NoError(t, world.purger.DeleteRemoteContent(ctx, wActorDID)) + retracted, err := world.votes.GetByATURI(ctx, pvVoteATURI) + require.NoError(t, err) + require.Equal(t, store.DeliveredStateUndone, retracted.DeliveredState, + "precondition: the purge retracted this vote — the decision is on the record") + + // The held settlement lands afterwards, which is the ordinary sequence: the + // hold is what made the purge see this vote at all. + world.settle(t) + + final, err := world.votes.GetByATURI(ctx, pvVoteATURI) + require.NoError(t, err, "the row still exists: a purge retracts a vote, it does not delete it") + assert.Equal(t, store.DeliveredStateUndone, final.DeliveredState, + "the vote must STILL be undone. The settlement is a true fact about the Like — the peer "+ + "accepted it — arriving after a decision that supersedes it: this actor is gone and "+ + "their votes are withdrawn. Writing 'delivered' over it re-counts a tombstoned "+ + "identity's vote in the fediverse-only tally 17b feeds, permanently, with no purge "+ + "left to re-run and no reseed that corrects it") +} + +// activitiesOfKind counts enqueued activities of one kind. +func activitiesOfKind(t *testing.T, conn *sql.DB, kind string) int { + t.Helper() + var n int + require.NoError(t, conn.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM outbound_activities WHERE kind = $1`, kind).Scan(&n)) + return n +} diff --git a/internal/outbound/worker.go b/internal/outbound/worker.go index bbc5137..8c499fa 100644 --- a/internal/outbound/worker.go +++ b/internal/outbound/worker.go @@ -9,6 +9,7 @@ import ( "log/slog" "net/http" "net/url" + "slices" "strings" "time" @@ -667,7 +668,12 @@ func (w *Worker) backoff(attempts int) time.Duration { // isRetraction reports whether a kind is a take-down that is exempt from the // consent recheck (a Delete or a vote Undo). -func isRetraction(kind string) bool { return kind == "Delete" || kind == "Undo" } +// +// The list is store.RetractionKinds, not a local copy: the QUEUE exempts exactly +// these kinds from a consent cancellation, and if the two definitions drifted, +// an activity this worker would still deliver could be cancelled out from under +// it — or the reverse. +func isRetraction(kind string) bool { return slices.Contains(store.RetractionKinds, kind) } // isDuplicate reports whether an HTTPError is Lemmy's duplicate-activity // response (a 400 whose body reports the activity was already received). diff --git a/internal/personas/serving.go b/internal/personas/serving.go index 62cdeb0..e30c7d8 100644 --- a/internal/personas/serving.go +++ b/internal/personas/serving.go @@ -234,8 +234,27 @@ func (s *Service) handleActorDocument(w http.ResponseWriter, r *http.Request, di // document keeps resolving, because every Note and Page already delivered // names it and revoking it would orphan the author reference on every // existing thread. + // + // THE BODY IS AN AS2 TOMBSTONE CARRYING THE PUBLIC KEY, and the key is there + // for a race this tier cannot otherwise win. The tombstone commits when the + // withdrawal is ENQUEUED; the worker POSTs it minutes later, with retries. In + // that window a peer verifying the signature on the very Delete that + // announces the withdrawal may re-dereference this actor, and a bare 410 + // leaves it with no key to verify with — so the erasure fails, retries and + // poisons, and the user is never actually withdrawn anywhere. + // + // Serving the key inside the Tombstone costs nothing a withdrawal cares + // about (the key was already public, and publicKey alone federates nothing) + // and gives a peer that reads the body what it needs to accept the last + // activity we will ever send as this actor. + // + // RESIDUAL, and it is real: a peer that reads only the STATUS still cannot + // verify. The complete fix is to keep serving the actor document until every + // delivery of the person-delete is terminal, which is a contract change — + // the tier's own test pins 410 immediately — and belongs to whoever changes + // that test. if actor.IsTombstoned() { - http.Error(w, "gone", http.StatusGone) + s.writeActorTombstone(w, actor) return } @@ -303,6 +322,13 @@ func (s *Service) handleOutbox(w http.ResponseWriter, r *http.Request, did strin http.NotFound(w, r) return } + // The outbox answers the same way the actor does. An actor that is Gone with + // a collection that is still 200 is a contradiction a peer has to resolve, + // and it invites exactly the re-fetch loop the 410 exists to end. + if actor.IsTombstoned() { + s.writeActorTombstone(w, actor) + return + } writeJSON(w, ap.ContentTypeActivityJSON, map[string]any{ "@context": asNamespace, "id": actor.ActorID + outboxSuffix, @@ -312,6 +338,29 @@ func (s *Service) handleOutbox(w http.ResponseWriter, r *http.Request, did strin }) } +// writeActorTombstone answers 410 Gone with an AS2 Tombstone for a withdrawn +// identity. formerType and deleted are what tell a peer WHAT is gone and WHEN, +// so it can retire its own copy rather than treat the status as a fetch failure. +func (s *Service) writeActorTombstone(w http.ResponseWriter, actor *store.APActor) { + w.Header().Set("Content-Type", ap.ContentTypeActivityJSON) + w.WriteHeader(http.StatusGone) + _ = json.NewEncoder(w).Encode(map[string]any{ + "@context": []any{asNamespace, securityNamespace}, + "id": actor.ActorID, + "type": ap.TypeTombstone, + "formerType": ap.TypePerson, + "deleted": actor.TombstonedAt.UTC().Format(time.RFC3339), + // See handleActorDocument: this is here so a peer can still verify the + // signature on the withdrawal activity itself, which may arrive after + // this document started answering Gone. + "publicKey": map[string]any{ + "id": actor.ActorID + "#main-key", + "owner": actor.ActorID, + "publicKeyPem": actor.PublicKeyPEM, + }, + }) +} + // handleWebFinger answers discovery for the local parts hosted on the ROUTED // Host. The lookup is (host, local part), so this origin can never answer for // an account it does not host, and two origins hosting the same local part diff --git a/internal/store/ap_actors.go b/internal/store/ap_actors.go index 1cb3207..0e22bfc 100644 --- a/internal/store/ap_actors.go +++ b/internal/store/ap_actors.go @@ -165,11 +165,25 @@ func (r *postgresAPActors) setEnabled(ctx context.Context, ex execer, did string // so the column answers "is this actor currently disabled, and since // when" rather than "was it ever disabled". enabled_at is not cleared on // disable: the last enable is history worth keeping. + // + // A TOMBSTONED ACTOR CAN NEVER BE RE-ENABLED, and the refusal is in the + // statement rather than in the callers because forgetting it is silent and + // unrecoverable. The destructive tier told every peer this identity was + // withdrawn and its document answers 410 forever; an enable arriving + // afterwards — enabled=true, or a DELETE of the opt-out record, which under + // default-on means the same thing — would resume signing new content as an + // actor peers were explicitly told is gone. "Irreversible" has to mean the + // identity cannot be brought back by writing a record. + // + // The row still MATCHES (so the caller gets a normal one-row result rather + // than a spurious NotFound); it is the VALUE that is forced. For an actor + // that was never tombstoned every branch below is exactly the old + // behaviour. query := ` UPDATE ap_actors SET - enabled = $2, - enabled_at = CASE WHEN $2 THEN now() ELSE enabled_at END, - disabled_at = CASE WHEN $2 THEN NULL ELSE now() END, + enabled = ($2 AND tombstoned_at IS NULL), + enabled_at = CASE WHEN ($2 AND tombstoned_at IS NULL) THEN now() ELSE enabled_at END, + disabled_at = CASE WHEN ($2 AND tombstoned_at IS NULL) THEN NULL ELSE now() END, updated_at = now() WHERE did = $1` return execOneRow(ctx, ex, "set enabled", did, query, did, enabled) diff --git a/internal/store/federation_prefs.go b/internal/store/federation_prefs.go index d3f3660..7722139 100644 --- a/internal/store/federation_prefs.go +++ b/internal/store/federation_prefs.go @@ -18,7 +18,7 @@ func NewFederationPrefs(db *sql.DB) FederationPrefs { return &postgresFederationPrefs{db: db} } -const federationPrefColumns = `did, enabled, delete_remote, source, updated_at` +const federationPrefColumns = `did, enabled, delete_remote, source, purged_at, updated_at` func (r *postgresFederationPrefs) Upsert(ctx context.Context, pref FederationPref) (*FederationPref, error) { if !pref.Source.Valid() { @@ -31,6 +31,13 @@ func (r *postgresFederationPrefs) Upsert(ctx context.Context, pref FederationPre // EVERY field is overwritten, delete_remote included. Re-enabling a user // must clear a previously requested deleteRemote: a stale destructive flag // sitting on an enabled row is a loaded gun pointed at task 17. + // + // EXCEPT purged_at, which is not in this statement at all — not in the + // INSERT, not in the DO UPDATE. It records that peers were actually asked to + // delete this user's content, which no preference write can make untrue, and + // leaving it out is what makes that impossible to undo by accident: a + // re-delivered account event upserting the same row cannot blank it, and no + // caller can set it by constructing a model. MarkPurged is the only writer. query := ` INSERT INTO federation_prefs (did, enabled, delete_remote, source, updated_at) VALUES ($1, $2, $3, $4, now()) @@ -69,19 +76,90 @@ func (r *postgresFederationPrefs) Delete(ctx context.Context, did string) error // Deleting a preference that never existed is success: an opt-out record // delete for a user who never opted out is the COMMON case, and absence is // exactly the state the delete is asking for. + // + // A COMMITTED PURGE IS NOT DELETABLE. Absence means default-on, so removing + // this row is how a user comes back — and a user whose content peers were + // already told to delete has nothing to come back to. The guard is in the + // statement because this delete is reachable from an ordinary record delete, + // 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`, did); err != nil { + `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) } return nil } +func (r *postgresFederationPrefs) MarkPurgedTx(ctx context.Context, tx *sql.Tx, did string) error { + if tx == nil { + return errors.NewValidationError("tx", "must not be nil") + } + return markPurged(ctx, tx, did) +} + +func (r *postgresFederationPrefs) MarkPurged(ctx context.Context, did string) error { + return markPurged(ctx, r.db, did) +} + +func markPurged(ctx context.Context, ex execer, did string) error { + // COALESCE: the FIRST commit is the one that happened. A retry that + // re-enqueues an idempotent withdrawal must not move the date, which is the + // only record of when this user's content was actually withdrawn. + result, err := ex.ExecContext(ctx, ` + UPDATE federation_prefs + SET purged_at = COALESCE(purged_at, now()), updated_at = now() + WHERE did = $1`, did) + if err != nil { + return fmt.Errorf("mark federation_pref %q purged: %w", did, err) + } + affected, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("mark federation_pref %q purged: rows affected: %w", did, err) + } + if affected == 0 { + return errors.NewNotFoundError("federation_pref", did) + } + return nil +} + +func (r *postgresFederationPrefs) ClearRequestedPurge(ctx context.Context, did string) (bool, error) { + // The predicate is the whole method, and both terms are load-bearing. + // + // source = 'account' — only a preference THIS tier wrote on the strength of + // a deletion event may be withdrawn by this tier. A user's own opt-out + // record says the same thing (enabled=false) and means something completely + // different: they asked. Clearing that would re-enable federation for + // somebody who never came back to ask for it. + // + // purged_at IS NULL — only a REQUEST may be withdrawn. Once peers have been + // asked to delete, the account being live again does not undo it, and + // resuming federation would publish under an identity those peers were told + // was gone. + result, err := r.db.ExecContext(ctx, ` + DELETE FROM federation_prefs + WHERE did = $1 AND source = $2 AND purged_at IS NULL`, + did, string(FederationPrefSourceAccount)) + if err != nil { + return false, fmt.Errorf("clear requested purge for %q: %w", did, err) + } + affected, err := result.RowsAffected() + if err != nil { + return false, fmt.Errorf("clear requested purge for %q: rows affected: %w", did, err) + } + return affected > 0, nil +} + func scanFederationPref(row rowScanner) (*FederationPref, error) { var pref FederationPref var source string - if err := row.Scan(&pref.DID, &pref.Enabled, &pref.DeleteRemote, &source, &pref.UpdatedAt); err != nil { + var purgedAt sql.NullTime + if err := row.Scan(&pref.DID, &pref.Enabled, &pref.DeleteRemote, &source, &purgedAt, &pref.UpdatedAt); err != nil { return nil, err } + if purgedAt.Valid { + pref.PurgedAt = &purgedAt.Time + } pref.Source = FederationPrefSource(source) return &pref, nil } diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index d7d85a1..18018e5 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -213,6 +213,11 @@ type APActors interface { // SetEnabled toggles federation for an actor: disabling stamps // disabled_at, re-enabling clears it and re-stamps enabled_at. A // missing actor is an error satisfying errors.IsNotFound. + // + // A TOMBSTONED actor is never re-enabled: the destructive tier is terminal, + // and the refusal lives in the statement so no caller can undo a withdrawal + // by writing a preference. Disabling one is still honoured (it is already + // disabled, and the write is idempotent). SetEnabled(ctx context.Context, did string, enabled bool) error // SetEnabledTx is SetEnabled on an existing transaction — the seam the @@ -516,11 +521,16 @@ type OutboundVotes interface { // from the one activity id. A miss is an error satisfying errors.IsNotFound. GetByActivityID(ctx context.Context, activityID string) (*OutboundVote, error) - // ListDeliveredForActor returns the actor's LIVE votes — the ones a peer - // still holds (delivered_state = 'delivered'). It is the destructive tier's - // input: a purged actor's standing votes are what the reseed keeps - // subtracting from a served score forever. - ListDeliveredForActor(ctx context.Context, actorDID string) ([]OutboundVote, error) + // ListStandingForActor returns the votes a PEER STILL HOLDS for this actor: + // those already settled as delivered, AND those whose delivery is held for + // settlement — accepted on the wire, with only our bookkeeping outstanding. + // + // The second half is what makes it correct as the destructive tier's input. + // A held vote reads 'pending' until the worker returns, and that return + // happens AFTER the withdrawal — so an enumeration of 'delivered' alone + // leaves a real vote standing on a real instance, attributed to an actor the + // bridge has told the world is gone, with nothing left that will notice. + ListStandingForActor(ctx context.Context, actorDID string) ([]OutboundVote, error) // SetDeliveredState transitions the delivery state. An unknown state is // an error satisfying errors.IsValidation; a missing vote is an error @@ -551,8 +561,32 @@ type FederationPrefs interface { Get(ctx context.Context, did string) (*FederationPref, error) // Delete removes the preference — the record-delete path, which restores - // the default-on state. Deleting a missing preference is a no-op success. + // the default-on state. Deleting a missing preference is a no-op success, + // and so is deleting one whose purge COMMITTED: absence means default-on, + // and a user whose content peers were already told to delete has nothing to + // come back to. Callers that report an outcome read the row back. Delete(ctx context.Context, did string) 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 + // is an error satisfying errors.IsNotFound, because a purge that committed + // with no preference to mark is a hole in the ordering, not a no-op. + MarkPurged(ctx context.Context, did string) error + + // MarkPurgedTx is MarkPurged on an existing transaction — the seam the + // purge itself uses so the marker lands atomically with the withdrawal it + // records. A nil tx is an error satisfying errors.IsValidation. + MarkPurgedTx(ctx context.Context, tx *sql.Tx, did string) error + + // ClearRequestedPurge withdraws a preference this tier wrote for a deletion + // that never happened — an account reported deleted, the purge failed, and + // the account is confirmed live again. It reports whether one was cleared. + // + // It clears ONLY (source=account AND purged_at IS NULL): a user's own + // opt-out is theirs to keep, and a committed purge is not undoable. Both + // terms live in the statement so a caller cannot reach past either. + ClearRequestedPurge(ctx context.Context, did string) (cleared bool, err error) } // Tombstones remembers AP object ids whose Delete arrived before (or @@ -689,11 +723,24 @@ type OutboundDeliveries interface { // posted. (A community BAN is the other shape and has its own statement.) CancelForActor(ctx context.Context, actorDID string) (int64, error) - // CancelForActorTx is CancelForActor on an existing transaction — the other - // half of the opt-out's atomic pair (see APActors.SetEnabledTx). A nil tx is + // CancelForActorTx is CancelForActor on an existing transaction. A nil tx is // an error satisfying errors.IsValidation. CancelForActorTx(ctx context.Context, tx *sql.Tx, actorDID string) (int64, error) + // CancelOutwardForActorTx is the CONSENT withdrawal, and the one an opt-out + // uses: it cancels everything that PUBLISHES for the actor and leaves their + // RETRACTIONS (Delete, Undo — see RetractionKinds) to go out. + // + // A user who deletes a post and then opts out must still have that delete + // delivered, or it stands on the peer forever; and the destructive tier's own + // Delete{Person} is a retraction too, so a sweeping cancel on a replay would + // silently cancel the erasure a previous attempt already committed. The + // worker draws the identical line at claim time — one list, so the two + // cannot disagree about what a stopped user is still owed. + // + // A nil tx is an error satisfying errors.IsValidation. + CancelOutwardForActorTx(ctx context.Context, tx *sql.Tx, actorDID string) (int64, error) + // CancelForCommunity moves every PENDING delivery on an ordering key (a // community AP id) to cancelled — a community deleted or unfollowed out // from under pending work. Returns how many rows were cancelled. diff --git a/internal/store/models.go b/internal/store/models.go index 7d0a991..99e4c81 100644 --- a/internal/store/models.go +++ b/internal/store/models.go @@ -231,12 +231,23 @@ const ( DeliveredStatePending DeliveredState = "pending" // DeliveredStateDelivered means a peer accepted the Like/Dislike. DeliveredStateDelivered DeliveredState = "delivered" - // DeliveredStateUndone is RESERVED and currently UNWRITTEN: task 15's worker - // DELETES the outbound_votes row on a successful Undo (clear-on-Undo) rather - // than transitioning it to "undone", so no code path ever sets this today. - // It is kept in the enum and the CHECK constraint (the migration is applied) - // against a future "keep the withdrawn-vote record" policy; Valid() still - // accepts it so a hand-set or legacy row round-trips. + // DeliveredStateUndone means the vote is NO LONGER LIVE on the peer as far as + // this bridge is concerned, so nothing may count it: the reseed subtracts + // only `delivered`, and the destructive tier enumerates only `delivered`. + // + // IT IS NO LONGER RESERVED, and its meaning is narrower than the obvious + // reading. Task 15's worker still DELETES the row on a successful Undo, so + // the ordinary retraction never passes through this state. ONE writer exists + // (task 17d's outbound.Purger): when an actor is withdrawn, their live votes + // are marked undone AT DECISION TIME, together with the Undo being enqueued + // — not when a peer confirms it. + // + // That distinction matters to anyone reasoning about the ledger: this value + // records OUR decision to stop counting a vote, not the peer's acceptance of + // the withdrawal. The two coincide for every path except a purge whose Undo + // never lands, where the peer may still hold a vote we have stopped counting + // — deliberately, because the alternative is subtracting forever on behalf + // of an identity that no longer exists. DeliveredStateUndone DeliveredState = "undone" ) @@ -357,7 +368,13 @@ type FederationPref struct { Enabled bool DeleteRemote bool Source FederationPrefSource - UpdatedAt time.Time + // PurgedAt is set once the destructive tier has actually asked peers to + // delete this user's content. It separates a withdrawal that was REQUESTED + // (nil — nothing irreversible has happened, so a live account may still be + // restored) from one that COMMITTED (set — there is nothing to come back + // to). Only MarkPurged writes it, and nothing clears it. + PurgedAt *time.Time + UpdatedAt time.Time } // InboxEvent is a received AP activity: the dedupe record AND the durable diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go index 66cfc46..099dbaf 100644 --- a/internal/store/outbound_deliveries.go +++ b/internal/store/outbound_deliveries.go @@ -8,6 +8,8 @@ import ( "strings" "time" + "github.com/lib/pq" + "tidepool/internal/errors" ) @@ -274,6 +276,13 @@ func (r *postgresOutboundDeliveries) CancelForActorTx(ctx context.Context, tx *s return cancelForActor(ctx, tx, actorDID) } +func (r *postgresOutboundDeliveries) CancelOutwardForActorTx(ctx context.Context, tx *sql.Tx, actorDID string) (int64, error) { + if tx == nil { + return 0, errors.NewValidationError("tx", "must not be nil") + } + return cancelOutwardForActor(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 @@ -319,11 +328,33 @@ const DeliveryHeldForSettlement = "ledger_unsettled" 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. -// Terminal deliveries are left untouched. Joined through -// outbound_activities.actor_did. +// RetractionKinds are the activity kinds that TAKE CONTENT DOWN: a Delete of a +// post or comment, and the Undo of a vote. They are exempt from every consent +// decision, in the queue and at the worker alike. +// +// The asymmetry is the point. A consent withdrawal means "stop publishing for +// me"; a retraction is the only way an opted-out user removes what is ALREADY +// published. Cancelling one leaves that content standing on the peer forever, +// which is the opposite of what was asked — the worker states exactly this and +// skips its consent recheck for these kinds, and a cancellation that swept them +// would undo that decision one layer up, where nothing observes it. +// +// It is ONE list shared with outbound.isRetraction so the queue and the worker +// cannot drift into disagreeing about which activities a stopped user is still +// owed. +var RetractionKinds = []string{"Delete", "Undo"} + +// notARetraction excludes those kinds from a cancellation. It reads from the +// joined outbound_activities row, so any statement using it must join a. +// cancelForActor is the SWEEPING cancel: every pending delivery of this actor's, +// retractions included, across every community. It is the kill switch and the +// operator's manual cancel — decisions that mean "stop the queue", not "stop +// publishing for this user". +// +// CancelOutwardForActor is the consent variant, and the difference between them +// is a user's ability to take their own content down. Two decisions, two +// statements: one predicate serving both is how the narrower decision silently +// acquires the wider one's reach. func cancelForActor(ctx context.Context, ex execer, actorDID string) (int64, error) { query := ` UPDATE outbound_deliveries d @@ -344,6 +375,46 @@ func cancelForActor(ctx context.Context, ex execer, actorDID string) (int64, err return affected, nil } +// cancelOutwardForActor is the CONSENT withdrawal: park the actor's pending +// OUTWARD work — everything that publishes — while leaving their retractions to +// go out. Never poisoned (this is not a failure), terminal rows untouched, held +// settlements untouched, across every community because the decision is about +// the actor. +// +// TWO THINGS DEPEND ON THE RETRACTION EXEMPTION, and the second is not obvious: +// +// 1. A user who deletes a post and then opts out must still have the delete +// delivered, or the post stays on Lemmy forever — the exact opposite of what +// opting out means. +// 2. THE DESTRUCTIVE TIER'S OWN WITHDRAWAL. Delete{Person} and the vote Undos +// are Deletes and Undos, enqueued by the purge on its own transaction. If a +// later replay of the opt-out record ran a sweeping cancel, it would cancel +// the erasure the previous attempt just committed — and nothing repairs it: +// the delivery insert returns the standing (cancelled) row by design, and +// the votes are already flipped, so the re-run enumerates nothing. Actor +// tombstoned, peers never told, "destructive opt-out applied" in the log. +func cancelOutwardForActor(ctx context.Context, ex execer, actorDID string) (int64, error) { + query := ` + UPDATE outbound_deliveries d + SET state = 'cancelled', claimed_until = NULL, updated_at = now() + FROM outbound_activities a + WHERE d.activity_id = a.activity_id + AND a.actor_did = $1 + AND d.state = 'pending'` + notHeldForSettlement + ` + AND a.kind <> ALL($2)` + + result, err := ex.ExecContext(ctx, query, actorDID, pq.Array(RetractionKinds)) + if err != nil { + return 0, fmt.Errorf("cancel outward outbound_deliveries for actor %q: %w", actorDID, err) + } + affected, err := result.RowsAffected() + if err != nil { + return 0, fmt.Errorf("cancel outward outbound_deliveries for actor %q: rows affected: %w", + actorDID, err) + } + return affected, nil +} + func (r *postgresOutboundDeliveries) CancelForCommunity(ctx context.Context, orderingKey string) (int64, error) { // A community deleted or unfollowed out from under pending work: park every // PENDING delivery on the ordering key as cancelled, leaving terminal rows. diff --git a/internal/store/outbound_standing_row_test.go b/internal/store/outbound_standing_row_test.go new file mode 100644 index 0000000..2e7ee38 --- /dev/null +++ b/internal/store/outbound_standing_row_test.go @@ -0,0 +1,237 @@ +package store + +import ( + "context" + "database/sql" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/testutil" +) + +// TASK 17d REVIEW — TWO RULES EVERY CANCELLATION AND EVERY REPLAY RESTS ON. +// +// (R1) THE STANDING ROW WINS. EnqueueTx is ON CONFLICT DO NOTHING plus a +// read-back, so re-enqueueing a delivery that already exists returns what is +// there rather than resetting it. That is deliberate and it is load-bearing in a +// direction nothing pins: the destructive tier commits on its own transaction +// while the rev gate commits later, so a failure in between REPLAYS the whole +// record — and every part of the purge is re-run. If a re-enqueue reset a +// delivered row to pending, the replay would re-POST a Delete{Person} the peer +// already applied; if it revived a cancelled one, a user's withdrawal would be +// un-withdrawn by a retry they never asked for. +// +// Today the rule is pinned only for pending↔pending, where nothing observable +// changes either way. These two tests pin it where the difference is real, so a +// future "fix" that resets the row reads as the regression it is instead of as +// an improvement. +// +// (R2) A DELIVERY HELD FOR SETTLEMENT IS NOT CANCELLABLE. notHeldForSettlement +// is pasted into FOUR statements and tested against one. The two below are the +// doors that had no test: 17c-3's community ban, and the community-wide sweep. +// The rule is the same everywhere — a cancellation answers "this must not go +// out", and a held delivery already went out, so cancelling it strands the +// ledger row it was held to settle — and the harm is remote from the code: +// a stranded row over-counts a served score forever AND hides the vote from the +// destructive tier's Undo enumeration. + +// standingTestDB is deliveryTestDB plus the ban table, which the ban door writes. +func standingTestDB(t *testing.T) *sql.DB { + t.Helper() + database := deliveryTestDB(t) + testutil.Truncate(t, database, "community_bans") + return database +} + +// seedOneDelivery inserts an activity and one pending delivery for it, and +// returns the delivery's key. +func seedOneDelivery(t *testing.T, database *sql.DB, activityID, actorDID, inbox, orderingKey string) { + t.Helper() + seedActivity(t, NewOutboundActivities(database), OutboundActivity{ + ActivityID: activityID, + ActorDID: actorDID, + Kind: "Create", + Payload: []byte(`{"type":"Create"}`), + }) + _, err := NewOutboundDeliveries(database).Enqueue(context.Background(), OutboundDelivery{ + ActivityID: activityID, TargetInbox: inbox, OrderingKey: orderingKey, + }) + require.NoError(t, err) +} + +// holdForSettlement drives a pending delivery into the held state through the +// real path — claim, then release under the settlement class — so the row is +// one the worker actually produces rather than one an UPDATE invented. +func holdForSettlement(t *testing.T, repo OutboundDeliveries, activityID, inbox string) { + t.Helper() + ctx := context.Background() + claimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + require.Equal(t, activityID, claimed.ActivityID, "the delivery under test is the one claimed") + require.Equal(t, inbox, claimed.TargetInbox) + _, applied, err := repo.Release(ctx, activityID, inbox, DeliveryHeldForSettlement, + "the peer accepted it; the local ledger write did not commit", 202, + time.Now().Add(time.Millisecond), *claimed.ClaimedUntil) + require.NoError(t, err) + require.True(t, applied, "the hold must actually be recorded") + + held, err := repo.Get(ctx, activityID, inbox) + require.NoError(t, err) + require.Equal(t, DeliveryStatePending, held.State, + "precondition: a held delivery is PENDING — that is what makes it re-claimable, and "+ + "also what puts it in reach of every cancellation") + require.Equal(t, DeliveryHeldForSettlement, held.LastErrorClass) +} + +// --------------------------------------------------------------------------- +// R1 — the standing row wins +// --------------------------------------------------------------------------- + +func TestOutboundDeliveries_EnqueueTxNeverRevivesATerminalDelivery(t *testing.T) { + for _, tc := range []struct { + name string + drive func(t *testing.T, repo OutboundDeliveries, activityID, inbox string) + want DeliveryState + why string + }{ + { + name: "cancelled", + drive: func(t *testing.T, repo OutboundDeliveries, _, _ string) { + cancelled, err := repo.CancelForActor(context.Background(), testDID) + require.NoError(t, err) + require.EqualValues(t, 1, cancelled) + }, + want: DeliveryStateCancelled, + why: "a cancelled delivery was withdrawn at the USER's request. A replay that " + + "revived it would publish on their behalf something they had already taken " + + "back — and the replay is not hypothetical: the destructive tier commits " + + "separately from the rev gate, so a rollback re-runs the whole enqueue", + }, + { + name: "delivered", + drive: func(t *testing.T, repo OutboundDeliveries, activityID, inbox string) { + ctx := context.Background() + claimed, err := repo.ClaimNext(ctx, time.Minute) + require.NoError(t, err) + _, applied, err := repo.MarkDelivered(ctx, activityID, inbox, 202, *claimed.ClaimedUntil) + require.NoError(t, err) + require.True(t, applied) + }, + want: DeliveryStateDelivered, + why: "a delivered row is the record that a peer HAS this activity. Resetting it to " + + "pending re-POSTs a Delete{Person} that instance already applied, and for the " + + "one tier that must never be asked twice", + }, + } { + t.Run(tc.name, func(t *testing.T) { + database := standingTestDB(t) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + seedOneDelivery(t, database, delActivityID, testDID, delTargetInbox, delOrderingKey) + tc.drive(t, repo, delActivityID, delTargetInbox) + + // The replay: the same activity enqueued to the same inbox again. + tx, err := database.BeginTx(ctx, nil) + require.NoError(t, err) + defer func() { _ = tx.Rollback() }() + returned, err := repo.EnqueueTx(ctx, tx, OutboundDelivery{ + ActivityID: delActivityID, TargetInbox: delTargetInbox, OrderingKey: delOrderingKey, + }) + require.NoError(t, err, "a re-enqueue is not an error: the caller's contract is a row that exists") + require.NoError(t, tx.Commit()) + + require.NotNil(t, returned) + assert.Equal(t, tc.want, returned.State, + "EnqueueTx must hand back the STANDING row, unchanged: %s", tc.why) + + stored, err := repo.Get(ctx, delActivityID, delTargetInbox) + require.NoError(t, err) + assert.Equal(t, tc.want, stored.State, "and the row in the table is unchanged too: %s", tc.why) + assert.Equal(t, 1, countDeliveries(t, database), + "with no second row beside it — the pair (activity, inbox) is one delivery, "+ + "however many times a replay asks for it") + }) + } +} + +func countDeliveries(t *testing.T, database *sql.DB) int { + t.Helper() + var n int + require.NoError(t, database.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM outbound_deliveries`).Scan(&n)) + return n +} + +// --------------------------------------------------------------------------- +// R2 — the other two doors into a held delivery +// --------------------------------------------------------------------------- + +// TestCommunityBans_BanLeavesADeliveryHeldForSettlement is the 17c-3 door. +// +// A ban cancels the banned author's pending work in that community, and a held +// delivery is pending. The vote it was holding to settle is one the community +// ALREADY has — banning the author does not un-cast it — so stranding the ledger +// row leaves that vote counted against the community's own score forever. +func TestCommunityBans_BanLeavesADeliveryHeldForSettlement(t *testing.T) { + database := standingTestDB(t) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + + // The held one, and an ordinary pending one beside it in the SAME community: + // without the second row, "the ban cancelled nothing" and "the ban spared the + // held row" are the same observation. + seedOneDelivery(t, database, delActivityID, testDID, delTargetInbox, delOrderingKey) + holdForSettlement(t, repo, delActivityID, delTargetInbox) + seedOneDelivery(t, database, delOtherActivityID, testDID, delTargetInbox+"/second", delOrderingKey) + + cancelled, err := NewCommunityBans(database).Ban(ctx, CommunityBan{ + CommunityDID: testCommunityDID, + SubjectDID: testDID, + CommunityAPID: delOrderingKey, + Reason: "spam", + }) + require.NoError(t, err) + assert.EqualValues(t, 1, cancelled, + "the ban cancels the author's ordinary queued work — the fixture is live") + + held, err := repo.Get(ctx, delActivityID, delTargetInbox) + require.NoError(t, err) + assert.Equal(t, DeliveryStatePending, held.State, + "but NOT the delivery held for settlement: the peer already has that activity, so "+ + "there is nothing left to stop — cancelling is terminal, the worker never returns, "+ + "and the ledger row it was held to settle is stranded where no reseed corrects it") + assert.Equal(t, DeliveryHeldForSettlement, held.LastErrorClass, + "and it keeps the class that tells the next claim to resume at the settlement") + + other, err := repo.Get(ctx, delOtherActivityID, delTargetInbox+"/second") + require.NoError(t, err) + assert.Equal(t, DeliveryStateCancelled, other.State, + "the ban still does its job on the work that has NOT gone out") +} + +// TestOutboundDeliveries_CancelForCommunityLeavesADeliveryHeldForSettlement is +// the community-wide door: a community deleted or unfollowed out from under +// pending work. Same rule, and the same reason — what has already reached the +// peer is not ours to un-send. +func TestOutboundDeliveries_CancelForCommunityLeavesADeliveryHeldForSettlement(t *testing.T) { + database := standingTestDB(t) + repo := NewOutboundDeliveries(database) + ctx := context.Background() + + seedOneDelivery(t, database, delActivityID, testDID, delTargetInbox, delOrderingKey) + holdForSettlement(t, repo, delActivityID, delTargetInbox) + seedOneDelivery(t, database, delOtherActivityID, testDID, delTargetInbox+"/second", delOrderingKey) + + cancelled, err := repo.CancelForCommunity(ctx, delOrderingKey) + require.NoError(t, err) + assert.EqualValues(t, 1, cancelled, "the ordinary pending delivery is cancelled — the fixture is live") + + held, err := repo.Get(ctx, delActivityID, delTargetInbox) + require.NoError(t, err) + assert.Equal(t, DeliveryStatePending, held.State, + "the held delivery survives a community-wide cancellation too: the settlement it is "+ + "waiting to finish is about an activity the community already received") +} diff --git a/internal/store/outbound_vote_terminal_test.go b/internal/store/outbound_vote_terminal_test.go new file mode 100644 index 0000000..543439e --- /dev/null +++ b/internal/store/outbound_vote_terminal_test.go @@ -0,0 +1,127 @@ +package store + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/errors" +) + +// TASK 17d — `undone` IS A DECISION, NOT A STAGE. +// +// Every other value in this column records where a vote GOT TO: pending is an +// intent, delivered is a peer's acceptance. `undone` records something else — +// OUR decision to stop counting a vote, written by the purge at the moment an +// actor is withdrawn, together with the Undo. It is the only value that is +// about an identity rather than about a message. +// +// So it is the only one that must not be overwritten by a later fact about the +// message. The route is real and it is the ordinary sequence, not a race: +// +// POST confirmed → ledger write fails → delivery HELD, vote 'pending' +// purge → enumerates the held vote, enqueues the Undo, marks 'undone' +// worker resumes → the held settlement lands → SetDeliveredState('delivered') +// +// and the last step un-retracts a vote for an actor this bridge has told the +// world is gone. Nothing revisits it: the purge is terminal and never re-runs, +// so 17b's reseed subtracts that vote from a served score forever, and a later +// operator reading the ledger sees a withdrawn user still holding live votes. +// +// THE GUARD BELONGS HERE, in the only writer, and not at its call sites: the +// settlement path cannot know it is racing a withdrawal, and any check it made +// would be a read outside the UPDATE's own snapshot. +// +// AND IT MUST SUCCEED. A guarded no-op is not a failure — the caller asked for +// something that is already decided, and the answer is "that is settled". If it +// came back as an error the worker would hold the delivery for settlement +// forever, retrying a write that can never apply, against a decision that will +// never change. That is why the missing-row case below is in the same file: the +// two must stay TELLABLE APART, and collapsing them is the tempting shortcut +// (`affected == 0 → nil`) that turns a real disagreement between the intent and +// the queue into a silence. + +// TestOutboundVotes_SetDeliveredStateCannotResurrectARetractedVote is the +// terminality itself. +func TestOutboundVotes_SetDeliveredStateCannotResurrectARetractedVote(t *testing.T) { + database := outboundTestDB(t) + repo := NewOutboundVotes(database) + ctx := context.Background() + + _, err := repo.Upsert(ctx, testOutboundVote()) + require.NoError(t, err) + require.NoError(t, repo.SetDeliveredState(ctx, testVoteATURI, DeliveredStateUndone), + "the purge's own write: the actor is withdrawn and this vote is retracted") + + // The late settlement, arriving exactly as the worker sends it. + err = repo.SetDeliveredState(ctx, testVoteATURI, DeliveredStateDelivered) + assert.NoError(t, err, + "a settlement that lands after a retraction is not an ERROR: the delivery it belongs "+ + "to succeeded, and the write it is asking for is simply already decided. Returning "+ + "an error here holds that delivery for settlement forever, retrying a write that "+ + "can never apply against a decision that will never change") + + got, err := repo.GetByATURI(ctx, testVoteATURI) + require.NoError(t, err) + require.NotNil(t, got) + assert.Equal(t, DeliveredStateUndone, got.DeliveredState, + "and the row STAYS undone: `undone` is this bridge's decision to stop counting the "+ + "vote of an actor it has withdrawn, not a stage the message passes through. "+ + "Flipping it back to delivered re-counts a vote for a tombstoned identity — the "+ + "reseed subtracts it from a served score forever, and nothing re-runs a purge") +} + +// TestOutboundVotes_SetDeliveredStateStillMovesALiveVote is the non-vacuity +// half, and it is not ceremony: the cheapest way to satisfy the test above is a +// guard that stops writing altogether. +func TestOutboundVotes_SetDeliveredStateStillMovesALiveVote(t *testing.T) { + database := outboundTestDB(t) + repo := NewOutboundVotes(database) + ctx := context.Background() + + _, err := repo.Upsert(ctx, testOutboundVote()) + require.NoError(t, err) + + require.NoError(t, repo.SetDeliveredState(ctx, testVoteATURI, DeliveredStateDelivered)) + got, err := repo.GetByATURI(ctx, testVoteATURI) + require.NoError(t, err) + assert.Equal(t, DeliveredStateDelivered, got.DeliveredState, + "an ordinary delivery success still flips pending -> delivered: the terminality is "+ + "about ONE value, and a guard that froze the column would take the vote ledger "+ + "— which 17b made an input to the number users read — permanently out of date") + + require.NoError(t, repo.SetDeliveredState(ctx, testVoteATURI, DeliveredStateUndone)) + got, err = repo.GetByATURI(ctx, testVoteATURI) + require.NoError(t, err) + assert.Equal(t, DeliveredStateUndone, got.DeliveredState, + "and delivered -> undone still applies: that transition IS the purge") +} + +// TestOutboundVotes_SetDeliveredStateOnAMissingVoteIsStillNotFound keeps the +// two nothings apart. +// +// A guarded no-op means "this is already decided". A missing row means the +// intent and the queue disagree about what exists — a delivery settling a vote +// nobody recorded. If the guard is written as "no rows changed, report success", +// the second becomes invisible, and the only signal that the two halves of the +// vote pipeline have diverged is gone. +func TestOutboundVotes_SetDeliveredStateOnAMissingVoteIsStillNotFound(t *testing.T) { + database := outboundTestDB(t) + repo := NewOutboundVotes(database) + ctx := context.Background() + + _, err := repo.Upsert(ctx, testOutboundVote()) + require.NoError(t, err) + require.NoError(t, repo.SetDeliveredState(ctx, testVoteATURI, DeliveredStateUndone), + "a retracted row exists in the table, so 'no row was updated' is genuinely ambiguous "+ + "unless the two cases are distinguished by more than the row count") + + err = repo.SetDeliveredState(ctx, testOtherVoteATURI, DeliveredStateDelivered) + require.Error(t, err, + "settling a vote we hold no state for is a real disagreement between the intent and "+ + "the delivery, and it must not be swallowed by the same branch that answers "+ + "'already retracted'") + assert.True(t, errors.IsNotFound(err), "want NotFound, got %v", err) +} diff --git a/internal/store/outbound_votes.go b/internal/store/outbound_votes.go index e676887..fc9792a 100644 --- a/internal/store/outbound_votes.go +++ b/internal/store/outbound_votes.go @@ -125,25 +125,38 @@ func (r *postgresOutboundVotes) GetByActivityID(ctx context.Context, activityID return vote, nil } -func (r *postgresOutboundVotes) ListDeliveredForActor(ctx context.Context, actorDID string) ([]OutboundVote, error) { +func (r *postgresOutboundVotes) ListStandingForActor(ctx context.Context, actorDID string) ([]OutboundVote, error) { if actorDID == "" { return nil, errors.NewValidationError("actor_did", "must not be empty") } - // LIVE means exactly delivered_state = 'delivered' — POSITIVE equality, per - // decision 16: those are the votes a peer still holds, and the same set the - // reseed subtracts from the origin's API tally. A purged actor leaving them - // standing is a number the reseed keeps subtracting from a score readers - // see, forever, on behalf of somebody who no longer exists. + // STANDING ON A PEER is a WIDER set than delivered_state = 'delivered', and + // the difference is a real vote on a real instance. // - // Served by the partial index on the same predicate (migration 029). + // A delivery HELD FOR SETTLEMENT has already been accepted by the peer — the + // POST returned, only our own bookkeeping failed — while its ledger row + // still reads 'pending' until the worker comes back to finish. Enumerating + // 'delivered' alone therefore misses a vote the peer demonstrably holds, and + // the miss is permanent rather than transient: that settlement lands AFTER + // the withdrawal meant to retract it, and a purge never re-runs. 17b's + // reseed then subtracts it from a served score forever, for somebody who no + // longer exists. + // + // The first term is served by the partial index (migration 029); the second + // is an EXISTS against the delivery whose id the vote already carries. query := `SELECT` + outboundVoteColumns + ` - FROM outbound_votes - WHERE actor_did = $1 AND delivered_state = 'delivered' - ORDER BY vote_at_uri` + FROM outbound_votes v + WHERE v.actor_did = $1 + AND (v.delivered_state = 'delivered' + OR EXISTS ( + SELECT 1 FROM outbound_deliveries d + WHERE d.activity_id = v.current_activity_id + AND d.state = 'pending' + AND d.last_error_class = '` + DeliveryHeldForSettlement + `')) + ORDER BY v.vote_at_uri` rows, err := r.db.QueryContext(ctx, query, actorDID) if err != nil { - return nil, fmt.Errorf("list delivered votes for %q: %w", actorDID, err) + return nil, fmt.Errorf("list standing votes for %q: %w", actorDID, err) } defer func() { _ = rows.Close() }() @@ -151,12 +164,12 @@ func (r *postgresOutboundVotes) ListDeliveredForActor(ctx context.Context, actor for rows.Next() { vote, err := scanOutboundVote(rows) if err != nil { - return nil, fmt.Errorf("scan delivered vote for %q: %w", actorDID, err) + return nil, fmt.Errorf("scan standing vote for %q: %w", actorDID, err) } votes = append(votes, *vote) } if err := rows.Err(); err != nil { - return nil, fmt.Errorf("list delivered votes for %q: %w", actorDID, err) + return nil, fmt.Errorf("list standing votes for %q: %w", actorDID, err) } return votes, nil } @@ -168,9 +181,21 @@ func (r *postgresOutboundVotes) SetDeliveredState(ctx context.Context, voteATURI if !state.Valid() { return errors.NewValidationError("delivered_state", "unknown state "+string(state)) } - result, err := r.db.ExecContext(ctx, - `UPDATE outbound_votes SET delivered_state = $2, updated_at = now() WHERE vote_at_uri = $1`, - voteATURI, string(state)) + // UNDONE IS TERMINAL HERE, and the ordering that makes this necessary is + // ordinary rather than exotic. A delivery HELD FOR SETTLEMENT has already + // been accepted by the peer, so a withdrawal can legitimately retract the + // vote while the worker is still on its way back to finish the bookkeeping; + // when it arrives it calls this method with `delivered`. Letting that late + // settlement win would re-establish exactly the state the erasure removed — + // the peer holds a vote we told them to drop, and 17b's reseed subtracts it + // from a served score forever. + // + // Re-setting undone stays allowed, so the write is idempotent. + result, err := r.db.ExecContext(ctx, ` + UPDATE outbound_votes SET delivered_state = $2, updated_at = now() + WHERE vote_at_uri = $1 + AND (delivered_state <> $3 OR $2 = $3)`, + voteATURI, string(state), string(DeliveredStateUndone)) if err != nil { return fmt.Errorf("set delivered_state for outbound_vote %q: %w", voteATURI, err) } @@ -179,10 +204,20 @@ func (r *postgresOutboundVotes) SetDeliveredState(ctx context.Context, voteATURI return fmt.Errorf("set delivered_state for outbound_vote %q: rows affected: %w", voteATURI, err) } if affected == 0 { - // Delivering a vote we hold no state for means the intent and the - // delivery disagree about what exists. That is a bug worth surfacing, - // not a no-op to swallow. - return errors.NewNotFoundError("outbound_vote", voteATURI) + // TWO DIFFERENT NOTHINGS, and they cannot share a branch. + // + // A row that is already retracted was deliberately not moved by the + // guard above: that is a DECIDED no-op and must report success. An error + // would fail the settlement, leaving the delivery held and retrying a + // write that can never apply, against a decision that will never change. + // + // A row that does not exist at all is the original finding this branch + // was written for — the intent and the delivery disagree about what + // exists — and is still worth surfacing. + if _, err := r.GetByATURI(ctx, voteATURI); err != nil { + return err + } + return nil } return nil } -- 2.51.2