diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index cc3a0fd..d1cf3e5 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -29,6 +29,7 @@ import ( "tidepool/internal/identity" "tidepool/internal/ingest" "tidepool/internal/materialize" + "tidepool/internal/optout" "tidepool/internal/outbound" "tidepool/internal/personas" "tidepool/internal/prune" @@ -737,12 +738,39 @@ func startConsumer( return nil, nil, fmt.Errorf("consumer: acceptance engine: %w", err) } + // The DESTRUCTIVE tier (decision 11's second tier): one Delete{Person, + // removeData:true} to every inbox this actor's content reached, an Undo for + // every vote peers still hold, and a 410 on the actor document. Reached from + // TWO doors — an explicit deleteRemote=true record, and a CONFIRMED account + // deletion — and never inferred from either. + purger := outbound.NewPurger(database, cfg.APUserOrigin, enqueuer).WithLogger(logger) + + // The TERMINAL tier (decision 19): a #account status of deleted is a claim + // about a moment that may have passed, so this confirms it against PLC and + // the PDS before anything irreversible is sent. The resolver is the + // confirmer — it already holds the guarded egress and the directory URL — + // and the purger only runs once that confirm comes back true. + terminator, err := optout.NewTerminator(optout.Options{ + Confirmer: resolver, + Prefs: store.NewFederationPrefs(database), + Deleter: purger, + Logger: logger, + }) + if err != nil { + return nil, nil, fmt.Errorf("consumer: account terminator: %w", err) + } + dispatcher, err := consume.NewDispatcher(consume.Options{ - DB: database, - Actors: minter, - Resolver: resolver, - Enqueuer: enqueuer, - Engine: engine, + DB: database, + Actors: minter, + Resolver: resolver, + Enqueuer: enqueuer, + Engine: engine, + Terminator: terminator, + // The record door: enabled=false + deleteRemote=true, written by the + // user themselves, so no confirmation is owed — the record IS the + // instruction. + RemoteDeleter: purger, // Reads committed records so a subject's community resolves for // mappings written before migration 016 filled community_did. Records: repoManager, diff --git a/internal/consume/account_status.go b/internal/consume/account_status.go new file mode 100644 index 0000000..65168da --- /dev/null +++ b/internal/consume/account_status.go @@ -0,0 +1,148 @@ +package consume + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "strings" +) + +// Account-status confirmation for decision 19's terminal tier. +// +// A #account event says the repo's status CHANGED; it does not prove what the +// status is now. The event may be stale (a reconnect rewind replays it), it may +// be a status the user has since reversed, and the action it triggers — asking +// every peer to delete the user's content — is one no peer undoes. So the +// terminal tier confirms against the identity's own sources before sending +// anything, and this is the read that does it. + +// atprotoPDSServiceID is the DID document service entry that names a repo's +// hosting 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 +// unbounded. +const maxRepoStatusBytes = 1 << 12 + +// AccountStatus reports whether the DID's repo is DELETED, confirmed against +// the identity's own sources: the PLC directory for the DID document, and the +// PDS that document names for the repo's status. +// +// THREE OUTCOMES, AND THEY MUST STAY THREE: +// +// (true, nil) — confirmed deleted. The destructive tier may run. +// (false, nil) — confirmed live. Nothing destructive; the event is handled. +// (_, err) — WE DO NOT KNOW. Nothing destructive, nothing recorded, retry. +// +// "We could not confirm" must never collapse into "confirmed not deleted", and +// it must never collapse into "confirmed deleted" either. Both directions are +// unrecoverable in opposite ways: one asks every peer to erase a user who never +// left, the other silently drops a real deletion and leaves their content +// federated forever. +// +// That collapse is 17c-3's P1-c one layer up. There, a ban's `expires` was read +// through a helper that returned nil for BOTH "absent" and "present but +// unparseable"; nil meant permanent; and an author was excluded forever because +// a timestamp did not parse. The shape of the bug was not the timestamp — it +// was two different facts arriving as one value on the field where the default +// was irreversible. This is that field, one layer up, so the error stays an +// error the whole way out. +// +// Every failure here is therefore an ERROR rather than a verdict: a directory +// that will not answer, a document with no PDS, a PDS that returns anything but +// a status we can read. Retrying costs a request; guessing costs a user's +// content. +func (r *HandleResolver) AccountStatus(ctx context.Context, did string) (bool, error) { + if err := validatePLCDID(did); err != nil { + return false, err + } + document, err := r.fetchDIDDocument(ctx, did) + if err != nil { + return false, fmt.Errorf("confirm account status for %s: %w", did, err) + } + endpoint, err := pdsEndpoint(document) + if err != nil { + // A document with no PDS is NOT read as "deleted", tempting as that is: + // the same shape appears while a document is mid-rewrite, and acting on + // it would erase a user whose hosting simply moved. + return false, fmt.Errorf("confirm account status for %s: %w", did, err) + } + return r.repoDeleted(ctx, endpoint, did) +} + +// didService is the sliver of a DID document's service list this reads. +type didService struct { + ID string `json:"id"` + Type string `json:"type"` + ServiceEndpoint string `json:"serviceEndpoint"` +} + +// 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) { + for _, service := range document.Service { + if !strings.HasSuffix(service.ID, 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 "", fmt.Errorf("DID document names no atproto PDS") +} + +// repoStatus is the sliver of com.atproto.sync.getRepoStatus this reads. +type repoStatus struct { + Active bool `json:"active"` + Status string `json:"status"` +} + +// repoStatusDeleted is the one status value that means the repo is gone. The +// others — takendown, suspended, deactivated — are states a user comes back +// from, which is the whole distinction decision 19 draws between the transient +// tier and this one. +const repoStatusDeleted = "deleted" + +// repoDeleted asks the PDS what it holds for this repo. +// +// ONLY an explicit, readable answer decides. A non-200 is unknown — a PDS that +// 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) + if err != nil { + return false, fmt.Errorf("build repo status request for %s: %w", did, err) + } + request.Header.Set("Accept", "application/json") + if r.userAgent != "" { + request.Header.Set("User-Agent", r.userAgent) + } + + response, err := r.httpClient.Do(request) + if err != nil { + return false, fmt.Errorf("fetch repo status for %s: %w", did, err) + } + defer func() { _ = response.Body.Close() }() + + if response.StatusCode != http.StatusOK { + return false, fmt.Errorf("fetch repo status for %s: pds returned %d", did, response.StatusCode) + } + var status repoStatus + if err := json.NewDecoder(io.LimitReader(response.Body, maxRepoStatusBytes)).Decode(&status); err != nil { + return false, fmt.Errorf("decode repo status for %s: %w", did, err) + } + // Both halves are required. `active` alone would read a suspension as a + // deletion, and `status` alone would believe a field the PDS may omit for a + // live repo. + return !status.Active && status.Status == repoStatusDeleted, nil +} diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index 8def8a1..0f932be 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -114,6 +114,30 @@ type PostIntent struct { // ActivityID reports the deterministic activity id. func (i PostIntent) ActivityID() string { return i.ID } +// PersonDeleteIntent withdraws a native user's whole identity from the +// fediverse: Delete{Person, removeData: true}, the destructive opt-out tier +// (decision 11's second tier, task 17d). +// +// It is ONE activity addressed to MANY inboxes — every instance this actor's +// content ever reached — which is what makes it the only intent here with no +// CommunityAPID: the targets come from the delivery history, not from a +// community, and the canonical payload is shared by every one of them. +// +// The shape lives here beside its siblings while the producer lives in +// outbound (the Purger), exactly as PostIntent's does. +// +// IRREVERSIBLE. Lemmy un-deletes a person on refetch; other software does not, +// and nothing in this codebase may promise resurrection. +type PersonDeleteIntent struct { + // ActorDID is the user being withdrawn. + ActorDID string + // ID is the deterministic activity id. + ID string +} + +// ActivityID reports the deterministic activity id. +func (i PersonDeleteIntent) ActivityID() string { return i.ID } + // OutboundEnqueuer is the task 15 seam. main wires the real persisting enqueuer // (outbound.Enqueuer, which writes outbound_activities/deliveries on the gate // tx) whenever CONSUMER_ENABLED; the noop is the consumer-disabled default. @@ -195,6 +219,9 @@ type Options struct { Votes store.OutboundVotes Communities store.Communities ObjectMappings store.APObjects + // Deliveries is the outbound queue the opt-out cancels an actor's pending + // work in — on the rev-gate transaction, together with the actor mirror. + Deliveries store.OutboundDeliveries // Bans reads whether the community a comment or vote is bound for has banned // its author (task 17c-3 review). The acceptance engine gates POSTS; nothing // else does, and a banned author's replies and votes enqueue to a community @@ -235,6 +262,7 @@ type Dispatcher struct { communities store.Communities moderation store.ObjectModeration bans store.CommunityBans + deliveries store.OutboundDeliveries records materialize.RecordGetter hosted *hostedRepos gate *RevGate @@ -306,6 +334,7 @@ func NewDispatcher(opts Options) (*Dispatcher, error) { communities: orDefault[store.Communities](opts.Communities, store.NewCommunities(opts.DB)), moderation: orDefault[store.ObjectModeration](opts.Moderation, store.NewObjectModeration(opts.DB)), bans: orDefault[store.CommunityBans](opts.Bans, store.NewCommunityBans(opts.DB)), + deliveries: orDefault[store.OutboundDeliveries](opts.Deliveries, store.NewOutboundDeliveries(opts.DB)), records: opts.Records, hosted: newHostedRepos(opts.DB), gate: NewRevGate(opts.DB), diff --git a/internal/consume/federation.go b/internal/consume/federation.go index 9922593..4c30da0 100644 --- a/internal/consume/federation.go +++ b/internal/consume/federation.go @@ -30,11 +30,18 @@ import ( // — because absence is what default-on looks like in this table. That also // clears any stored deleteRemote: a stale destructive flag on a re-enabled // user is a loaded gun pointed at the task 17 tier. -func (d *Dispatcher) handleFederation(ctx context.Context, _ *sql.Tx, did string, commit *CommitEvent) error { +// +// tx is the REV-GATE's transaction, and the disable path writes on it: the +// actor mirror and the cancellation of that actor's queued deliveries are one +// decision, and they commit with the gate advance or not at all. Atomicity here +// is structural rather than defended — there is no window in which one landed +// and the other did not, and a failure leaves the gate un-advanced so the record +// replays and re-applies both. +func (d *Dispatcher) handleFederation(ctx context.Context, tx *sql.Tx, did string, commit *CommitEvent) error { if commit.Operation == operationDelete { // A delete commit carries no record body, which costs nothing here: // the DID is the whole question. - return d.restoreDefaultFederation(ctx, did) + return d.restoreDefaultFederation(ctx, tx, did) } enabled, ok := boolField(commit.Record, "enabled") @@ -47,16 +54,19 @@ func (d *Dispatcher) handleFederation(ctx context.Context, _ *sql.Tx, did string ErrPermanentEvent, did) } if enabled { - return d.restoreDefaultFederation(ctx, did) + return d.restoreDefaultFederation(ctx, tx, did) } // deleteRemote is optional and defaults to false: the soft tier. Nothing // destructive is ever INFERRED — only an explicit true escalates. deleteRemote, _ := boolField(commit.Record, "deleteRemote") - // The preference is recorded BEFORE the destructive seam is reached. - // Peers that honor a Delete cannot restore what they dropped, so the - // user's intent must survive a crash between recording and sending. + // The preference is recorded BEFORE the destructive seam is reached, and + // deliberately NOT on the transaction below. Peers that honor a Delete + // cannot restore what they dropped, so the user's intent must be durable + // before anything is sent — and the gate transaction has not committed by + // the time the destructive seam runs. It is the authority besides: it + // answers for every DID, including the ones with no actor to mirror onto. if _, err := d.prefs.Upsert(ctx, store.FederationPref{ DID: did, Enabled: false, @@ -66,13 +76,56 @@ func (d *Dispatcher) handleFederation(ctx context.Context, _ *sql.Tx, did string return fmt.Errorf("record federation opt-out for %s: %w", did, err) } - if err := d.mirrorActorEnabled(ctx, did, false); err != nil { - return err + // 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) + if err != nil { + return fmt.Errorf("cancel queued deliveries for %s: %w", did, err) + } + if cancelled > 0 { + d.logger.Info("federation opt-out cancelled queued deliveries", + slog.String("did", did), slog.Int64("cancelled", cancelled)) } - if !deleteRemote { - return nil + if deleteRemote { + if err := d.deleteRemoteContent(ctx, did); err != nil { + return err + } } + + // The mirror answers the first fact — the consumer, the admission gate and + // the delivery claim all read it — and it lands LAST, for a reason that is + // about locks rather than meaning: it UPDATEs the actor row, and the + // destructive tier updates that same row from its own transaction. Holding + // this one across that call deadlocks the two against each other, and the + // handler cannot release a lock it holds without committing. + // + // It still rides the rev-gate transaction, so the cancellation and the flag + // commit together with the gate advance: nothing retries the missing half, + // because the record is applied ONCE under a gate that rejects the replay. + return d.mirrorActorEnabled(ctx, tx, did, false) +} + +// deleteRemoteContent runs the DESTRUCTIVE tier: peers are asked to delete what +// they already hold, the votes they still count are retracted, and the identity +// stops resolving. +// +// It runs on its OWN transaction rather than the gate's (the seam takes no tx), +// which is why the caller must not be holding a lock on anything it writes. The +// consequence of that split is stated plainly: if the gate transaction later +// rolls back, the withdrawal has still been enqueued and the record replays — +// which is safe only because every part of the purge is idempotent (a +// deterministic activity id, a delivery insert that returns the standing row, +// a tombstone that keeps its first timestamp). +// +// IRREVERSIBLE. Peers that honour a Delete cannot restore what they dropped, +// and Lemmy's own un-delete on refetch is not something other software promises. +func (d *Dispatcher) deleteRemoteContent(ctx context.Context, did string) error { if d.remoteDeleter == nil { // A deployment where the destructive tier has not landed yet. The // request is already recorded, so it can be acted on from the backlog @@ -93,19 +146,24 @@ func (d *Dispatcher) handleFederation(ctx context.Context, _ *sql.Tx, did string // restoreDefaultFederation is the re-enable path shared by enabled=true and a // record delete: the preference row goes away (absence IS default-on) and an // EXISTING actor is re-enabled under its original identity. -func (d *Dispatcher) restoreDefaultFederation(ctx context.Context, did string) error { +// +// Cancelled deliveries are NOT resurrected. Re-enabling restores the user's +// 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. +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, did, true) + return d.mirrorActorEnabled(ctx, tx, did, true) } // mirrorActorEnabled updates the ap_actors mirror IF the actor exists. A // missing actor is a no-op success, not an error: an opt-out (or an enable) // from a DID that has never federated anything is the ordinary case, and // actors mint at the first federating interaction — never here. -func (d *Dispatcher) mirrorActorEnabled(ctx context.Context, did string, enabled bool) error { - err := d.apActors.SetEnabled(ctx, did, enabled) +func (d *Dispatcher) mirrorActorEnabled(ctx context.Context, tx *sql.Tx, did string, enabled bool) error { + err := d.apActors.SetEnabledTx(ctx, tx, did, enabled) if errors.IsNotFound(err) { d.logger.Debug("federation preference for a DID with no actor", slog.String("did", did), slog.Bool("enabled", enabled)) diff --git a/internal/consume/resolver.go b/internal/consume/resolver.go index b49c989..5d04d68 100644 --- a/internal/consume/resolver.go +++ b/internal/consume/resolver.go @@ -172,9 +172,12 @@ func (r *HandleResolver) ResolveDIDHandle(ctx context.Context, did string) (stri return handle, nil } -// didDocument is the sliver of a DID document this resolver reads. +// didDocument is the sliver of a DID document this resolver reads: the handle +// claims, and the services that name where the repo is hosted (see +// account_status.go, which confirms a deletion against that PDS). type didDocument struct { - AlsoKnownAs []string `json:"alsoKnownAs"` + AlsoKnownAs []string `json:"alsoKnownAs"` + Service []didService `json:"service"` } // fetchDIDDocument reads the DID document from the PLC directory. Every diff --git a/internal/consume/votes.go b/internal/consume/votes.go index 7d6ce22..1ba78a4 100644 --- a/internal/consume/votes.go +++ b/internal/consume/votes.go @@ -208,7 +208,16 @@ func (d *Dispatcher) applyVoteDelete(ctx context.Context, tx *sql.Tx, did string // operationUndo is not a commit operation: it is the outbound op a vote delete // becomes, kept distinct from "delete" because an Undo names the activity it // withdraws rather than an object. -const operationUndo = "undo" +// +// It is EXPORTED as OperationUndo because the activity id derives from it, and +// a second producer of vote Undos now exists: the destructive opt-out tier +// retracts a withdrawn actor's live votes (outbound.Purger). Two copies of the +// string would mint two different ids for the same operation. +const operationUndo = OperationUndo + +// OperationUndo is the outbound op a vote retraction is derived under. See +// operationUndo. +const OperationUndo = "undo" // communityAPID resolves a community's AP Group id for addressing. // diff --git a/internal/db/migrations/028_federation_pref_account_source.sql b/internal/db/migrations/028_federation_pref_account_source.sql new file mode 100644 index 0000000..41bcd8a --- /dev/null +++ b/internal/db/migrations/028_federation_pref_account_source.sql @@ -0,0 +1,31 @@ +-- +goose Up +-- Task 17d: a third provenance for a federation preference. +-- +-- federation_prefs.source answers "how do we know this?", and the column exists +-- precisely because the answers have different staleness (see migration 018). +-- Until now there were two: a social.coves.bridge.federation RECORD the consumer +-- saw, and a PROBE where the bridge fetched that record itself. +-- +-- The terminal tier (decision 19) writes a third. When a #account event says a +-- repo is deleted, the bridge CONFIRMS it against the identity's own sources — +-- the PLC directory, then the PDS that document names — and records the result. +-- That is not a federation record at all: the user never wrote one, and there is +-- no record left to fetch. Filing it under 'probe' would claim we read a +-- preference they expressed, when what we read is that their account is gone — +-- and this column is the one place an operator can tell a user who OPTED OUT +-- from a user who was DELETED, which is the difference between a decision they +-- can reverse and one they cannot. +ALTER TABLE federation_prefs DROP CONSTRAINT IF EXISTS federation_prefs_source_check; +ALTER TABLE federation_prefs ADD CONSTRAINT federation_prefs_source_check + CHECK (source IN ('record', 'probe', 'account')); + +-- +goose Down +-- Rows written by the terminal tier are rewritten to 'probe' rather than +-- deleted: the PREFERENCE is what stops the bridge federating for a deleted +-- account, so dropping the rows to satisfy a narrower constraint would resume +-- federating on their behalf. Losing the provenance is recoverable; losing the +-- preference is not. +UPDATE federation_prefs SET source = 'probe' WHERE source = 'account'; +ALTER TABLE federation_prefs DROP CONSTRAINT IF EXISTS federation_prefs_source_check; +ALTER TABLE federation_prefs ADD CONSTRAINT federation_prefs_source_check + CHECK (source IN ('record', 'probe')); diff --git a/internal/db/migrations/029_actor_tombstone.sql b/internal/db/migrations/029_actor_tombstone.sql new file mode 100644 index 0000000..7d8d7e5 --- /dev/null +++ b/internal/db/migrations/029_actor_tombstone.sql @@ -0,0 +1,41 @@ +-- +goose Up +-- Task 17d: the DESTRUCTIVE tier's two reads. +-- +-- 1. tombstoned_at — the actor document stops resolving. +-- +-- ap_actors already had three lifecycle columns and none of them can express +-- this. `enabled` gates DISCOVERY: a disabled actor's local part stops resolving +-- through webfinger while its document keeps being served, deliberately, because +-- every Note and Page the bridge already delivered names that actor and a +-- dangling reference orphans the thread it hangs in. `delivery_paused` is about +-- outbound traffic and has no read-side meaning at all. +-- +-- A withdrawal is the opposite decision, and it is the ONLY difference a peer +-- can observe between the two tiers: the document answers 410 GONE. Not 404 — +-- "never heard of them" reads as a lookup failure and invites the peer to try +-- again, while Gone is the statement that lets it stop asking and clean up. +-- +-- It is a TIMESTAMP rather than a flag because when a withdrawal happened is the +-- question anyone asks afterwards, and because IRREVERSIBILITY is the point: no +-- code path clears this column. Peers that honour a Delete cannot restore what +-- they dropped, so nothing here promises resurrection. +ALTER TABLE ap_actors ADD COLUMN tombstoned_at TIMESTAMPTZ; + +-- 2. The purge's vote query. +-- +-- A purged actor's LIVE votes have to be retracted, and "live" is exactly +-- delivered_state = 'delivered' (task 17b: the reseed subtracts precisely those +-- from the origin's API tally, so a vote left standing for an actor that no +-- longer exists is a number the reseed keeps subtracting from a served score +-- forever). The lookup is by actor, which no index served: outbound_votes is +-- indexed for the subject-side reseed, not the actor-side purge. +-- +-- PARTIAL on the same predicate the reader uses, so the index holds only the +-- rows a purge can act on — 17b's ruling that the subject index be partial, one +-- axis over. +CREATE INDEX outbound_votes_actor_delivered_idx + ON outbound_votes (actor_did) WHERE delivered_state = 'delivered'; + +-- +goose Down +DROP INDEX IF EXISTS outbound_votes_actor_delivered_idx; +ALTER TABLE ap_actors DROP COLUMN IF EXISTS tombstoned_at; diff --git a/internal/ingest/ingest_test.go b/internal/ingest/ingest_test.go index 7d0a357..8bfe6db 100644 --- a/internal/ingest/ingest_test.go +++ b/internal/ingest/ingest_test.go @@ -11,6 +11,7 @@ import ( "log/slog" "net/http" "net/http/httptest" + "net/url" "os" "path/filepath" "strings" @@ -687,6 +688,15 @@ func (h *harness) subscribeTechnology() *remoteActor { return group } +// originOf is the scheme://host of an AP id. +func originOf(apID string) string { + parsed, err := url.Parse(apID) + if err != nil || parsed.Host == "" { + return apID + } + return parsed.Scheme + "://" + parsed.Host +} + // subscribeCommunityURL subscribes to a SECOND community through the real admin // path — resolve, mint, Follow, Accept — and returns its signing handle. // @@ -706,8 +716,13 @@ func (h *harness) subscribeCommunityURL(apGroupID, username string) *remoteActor "id": apGroupID, "preferredUsername": username, "inbox": apGroupID + "/inbox", - "endpoints": map[string]any{"sharedInbox": "https://lemmy.world/inbox"}, - "published": "2024-01-01T00:00:00.000000Z", + // The shared inbox is derived from the community's OWN host, so + // co-hosted communities share one (which is what makes ordering_key + // rather than target_inbox the thing that carries scope) while + // communities on different instances do not (which is what makes a + // fan-out across instances expressible at all). + "endpoints": map[string]any{"sharedInbox": originOf(apGroupID) + "/inbox"}, + "published": "2024-01-01T00:00:00.000000Z", }) rec := h.adminRequest(http.MethodPost, "/admin/communities", diff --git a/internal/ingest/moderation_terminal_test.go b/internal/ingest/moderation_terminal_test.go index 2659d99..a4c1cb7 100644 --- a/internal/ingest/moderation_terminal_test.go +++ b/internal/ingest/moderation_terminal_test.go @@ -150,12 +150,17 @@ func newModerationWorld(t *testing.T, h *harness) moderationWorld { }) require.NoError(t, err) dispatcher, err := consume.NewDispatcher(consume.Options{ - DB: h.db, - Actors: userOrigin, - Enqueuer: enqueuer, - Resolver: mtResolver{}, - Engine: engine, - UserOrigin: mtUserOrigin, + DB: h.db, + Actors: userOrigin, + Enqueuer: enqueuer, + Resolver: mtResolver{}, + Engine: engine, + // The DESTRUCTIVE opt-out tier, wired exactly as production wires it. + // A nil deleter makes deleteRemote=true a Warn and a no-op, which would + // let every destructive-tier assertion pass against a tier that was + // never reached. + RemoteDeleter: outbound.NewPurger(h.db, mtUserOrigin, enqueuer), + UserOrigin: mtUserOrigin, }) require.NoError(t, err) diff --git a/internal/ingest/optout_destructive_test.go b/internal/ingest/optout_destructive_test.go new file mode 100644 index 0000000..46bfd25 --- /dev/null +++ b/internal/ingest/optout_destructive_test.go @@ -0,0 +1,237 @@ +package ingest + +import ( + "context" + "database/sql" + "net/http" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TASK 17d — THE DESTRUCTIVE TIER, AND WHAT SEPARATES IT FROM THE SOFT ONE. +// +// The soft tier stops the bridge SPEAKING for someone. The destructive tier +// withdraws what it already said. A user reaches it only by writing +// deleteRemote=true — nothing infers it — because peers that honour a Delete +// cannot restore what they dropped. +// +// Three things have to happen, and each is unreachable from the others: +// +// 1. Delete{Person, removeData:true} to EVERY community inbox this actor's +// content ever reached. One inbox is not "most of it": the instances that +// do not receive it keep serving the user's posts forever, and there is no +// second attempt — asking twice risks deleting content a re-enabled user +// has since restored. +// 2. An Undo for every vote those peers still hold. A purged actor leaving +// standing tallies is not cosmetic: Tidepool's aggregate is the +// fediverse-only tally, and the reseed SUBTRACTS live delivered outbound +// votes from the API count. A vote nobody can attribute to a living actor +// is a number the reseed keeps subtracting against, forever, on a subject +// whose score is served to readers. +// 3. The actor document stops resolving — 410 Gone, not 404. Gone is the +// statement "this existed and was withdrawn", which is what a peer needs to +// stop retrying and clean up; 404 reads as "never heard of them", which +// several implementations treat as a transient lookup failure. +// +// THE PAIRING IS THE POINT. The soft tier keeps serving the document, because +// every Note and Page already delivered names this actor and a broken author +// reference orphans every existing thread. The destructive tier revokes it. Both +// assertions live in this file, on two actors, in one run — the difference +// between the tiers is a status code, and prose is not where that belongs. + +const ( + odDestructiveRev = "3lzodrev000100" + odSoftRev = "3lzodrev000101" + + odPurgeInARKey = "3lzodpurge0001" + odPurgeInBRKey = "3lzodpurge0002" + odSoftRKey = "3lzodpurge0003" + + // A community on a DIFFERENT INSTANCE. The fixture's two communities are + // co-hosted and therefore share one inbox — correct for the ban work, and + // useless here: a fan-out that reaches "every inbox" is indistinguishable + // from one that reaches the first when there is only one inbox to reach. + // An erasure request is about INSTANCES, so the fixture needs two. + odFarCommunity = "https://lemmy.zip/c/purgetest" + odFarName = "purgetest" + odFarRKey = "3lzodpurge0004" +) + +// TestTheDestructiveTierWithdrawsEverythingTheSoftTierKeeps is the outer +// contract for deleteRemote=true. +func TestTheDestructiveTierWithdrawsEverythingTheSoftTierKeeps(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + // --- GIVEN: an author whose content reached TWO instances, holding one + // LIVE vote and one already-retracted one. + admitPost(t, world, mtAuthorDID, odPurgeInARKey, world.communityADID, "3lzodrev000110", 1_775_000_030_000_001) + admitPost(t, world, mtAuthorDID, odPurgeInBRKey, world.communityBDID, "3lzodrev000111", 1_775_000_030_000_002) + + // ...and one on another instance, which is where "every inbox" starts to + // mean something. + h.subscribeCommunityURL(odFarCommunity, odFarName) + admitPost(t, world, mtAuthorDID, odFarRKey, testDIDFor(odFarName, "lemmy.zip"), + "3lzodrev000113", 1_775_000_030_000_004) + + inboxes := deliveryInboxesFor(t, h.db, mtAuthorDID) + require.Len(t, inboxes, 2, + "precondition: this actor's content reached two distinct inboxes — one is a fan-out "+ + "that cannot be told from a single delivery") + + // A second actor takes the SOFT tier in the same run, so the two outcomes + // are compared under one set of conditions. + admitPost(t, world, mtCommenterDID, odSoftRKey, world.communityADID, "3lzodrev000112", 1_775_000_030_000_003) + + liveVote := seedDeliveredVote(t, h.db, mtAuthorDID, mtPostATURI, world.communityADID, "delivered") + deadVote := seedDeliveredVote(t, h.db, mtAuthorDID, odVoteSubjectATURI, world.communityADID, "undone") + + activitiesBefore := rowCount(t, h.db, "outbound_activities") + + // --- WHEN: the destructive record, and the soft one beside it. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, mtAuthorDID, odDestructiveRev, "create", false, true, 1_775_000_031_000_001))) + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, mtCommenterDID, odSoftRev, "create", false, false, 1_775_000_031_000_002))) + + // --- THEN (1): one Delete{Person}, delivered to every inbox they reached. + deleteActivity := activityOfKind(t, h.db, mtAuthorDID, "Delete") + require.NotEmpty(t, deleteActivity, + "a Delete{Person} must be enqueued: the soft tier tells peers nothing, and this tier "+ + "is defined by telling them — a destructive opt-out that sends no activity is a "+ + "soft opt-out the user did not choose") + assert.ElementsMatch(t, inboxes, deliveryInboxesForActivity(t, h.db, deleteActivity), + "and it reaches EVERY inbox this actor's content ever reached: the instances it "+ + "misses keep serving their posts, and no second attempt is coming — this is the "+ + "fan-out the enqueuer's idempotency shortcut silently truncates to one") + + // --- THEN (2): an Undo for the vote peers still hold, and only that one. + assert.Equal(t, 1, undoActivitiesFor(t, h.db, mtAuthorDID), + "exactly ONE Undo: for the vote still standing on a peer, and NOT for the one already "+ + "retracted. A purged actor's live vote is a number the reseed keeps subtracting "+ + "from a served score forever; a second Undo for a vote nobody holds is an activity "+ + "the peer cannot match to anything") + assert.Equal(t, "undone", voteState(t, h.db, liveVote), + "and the live vote's own state says so afterwards, or the next reseed subtracts it again") + assert.Equal(t, "undone", voteState(t, h.db, deadVote), "while the retracted one is unchanged") + + assert.Greater(t, rowCount(t, h.db, "outbound_activities"), activitiesBefore, + "the tier really did enqueue: the assertions above must not be satisfied by silence") + + // --- THEN (3): the actor document is GONE, and the soft actor's is not. + assert.Equal(t, http.StatusGone, actorDocStatus(t, h, mtAuthorDID), + "the purged actor's document returns 410 GONE: the user asked to be withdrawn, and a "+ + "document that still resolves invites peers to keep re-fetching an identity we "+ + "promised to retract. 404 is the wrong answer — 'never heard of them' reads as a "+ + "lookup failure, while Gone is the statement that lets a peer stop asking") + + assert.Equal(t, http.StatusOK, actorDocStatus(t, h, mtCommenterDID), + "while the SOFT opt-out's document still resolves: every Note and Page already "+ + "delivered names that actor, and revoking it orphans the author reference on every "+ + "existing thread. This is the whole difference between the tiers, and it is one "+ + "status code") +} + +// odVoteSubjectATURI is a second subject to hang the already-retracted vote on. +const odVoteSubjectATURI = "at://" + mtAuthorDID + "/social.coves.community.postv2/3lzodvotesub01" + +// seedDeliveredVote writes the outbound vote state a purge has to act on. The +// delivered/undone distinction is the whole question — "live" means a peer still +// holds it — and driving a real delivery would test the worker instead. +func seedDeliveredVote(t *testing.T, db *sql.DB, actorDID, subjectATURI, communityDID, state string) string { + t.Helper() + voteATURI := "at://" + actorDID + "/social.coves.feed.vote/" + state + "-vote" + _, err := db.ExecContext(context.Background(), ` + INSERT INTO outbound_votes ( + vote_at_uri, actor_did, subject_at_uri, subject_ap_id, community_did, + direction, current_activity_id, delivered_state) + VALUES ($1, $2, $3, $4, $5, 'up', $6, $7)`, + voteATURI, actorDID, subjectATURI, mtUserOrigin+"/ap/object/"+subjectATURI, + communityDID, mtUserOrigin+"/ap/activity/"+state+"-vote", state) + require.NoError(t, err, "seed a %s vote for %s", state, actorDID) + return voteATURI +} + +func voteState(t *testing.T, db *sql.DB, voteATURI string) string { + t.Helper() + var state string + require.NoError(t, db.QueryRowContext(context.Background(), + `SELECT delivered_state FROM outbound_votes WHERE vote_at_uri = $1`, voteATURI).Scan(&state)) + return state +} + +// deliveryInboxesFor lists the distinct inboxes an actor's content has reached — +// the delivery history the fan-out is built from. No production query answers +// this today; that is new surface the tier needs. +func deliveryInboxesFor(t *testing.T, db *sql.DB, actorDID string) []string { + t.Helper() + return queryStrings(t, db, ` + SELECT DISTINCT d.target_inbox + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE a.actor_did = $1`, actorDID) +} + +func deliveryInboxesForActivity(t *testing.T, db *sql.DB, activityID string) []string { + t.Helper() + return queryStrings(t, + db, `SELECT target_inbox FROM outbound_deliveries WHERE activity_id = $1`, activityID) +} + +// activityOfKind returns one activity id of the given kind for an actor, or "". +func activityOfKind(t *testing.T, db *sql.DB, actorDID, kind string) string { + t.Helper() + ids := queryStrings(t, db, + `SELECT activity_id FROM outbound_activities WHERE actor_did = $1 AND kind = $2`, + actorDID, kind) + if len(ids) == 0 { + return "" + } + require.Len(t, ids, 1, + "one %s activity for %s: the erasure is ONE request fanned out, not one per peer", + kind, actorDID) + return ids[0] +} + +func undoActivitiesFor(t *testing.T, db *sql.DB, actorDID string) int { + t.Helper() + var n int + require.NoError(t, db.QueryRowContext(context.Background(), + `SELECT COUNT(*) FROM outbound_activities WHERE actor_did = $1 AND kind = 'Undo'`, + actorDID).Scan(&n)) + return n +} + +// actorDocStatus fetches the persona's actor document from the origin the +// bridge really serves, and reports the status. It goes over HTTP on purpose: +// what a peer receives is the contract, and a store flag that no handler reads +// would satisfy any assertion made against the database instead. +func actorDocStatus(t *testing.T, h *harness, did string) int { + t.Helper() + req, err := http.NewRequestWithContext(context.Background(), http.MethodGet, + "http://"+h.fixtures.Listener.Addr().String()+"/ap/actor/"+did, nil) + require.NoError(t, err) + req.Host = "coves.social" + resp, err := http.DefaultClient.Do(req) + require.NoError(t, err) + defer func() { require.NoError(t, resp.Body.Close()) }() + return resp.StatusCode +} + +func queryStrings(t *testing.T, db *sql.DB, query string, args ...any) []string { + t.Helper() + rows, err := db.QueryContext(context.Background(), query, args...) + require.NoError(t, err) + defer func() { require.NoError(t, rows.Close()) }() + var out []string + for rows.Next() { + var value string + require.NoError(t, rows.Scan(&value)) + out = append(out, value) + } + require.NoError(t, rows.Err()) + return out +} diff --git a/internal/ingest/optout_test.go b/internal/ingest/optout_test.go new file mode 100644 index 0000000..255bf41 --- /dev/null +++ b/internal/ingest/optout_test.go @@ -0,0 +1,171 @@ +package ingest + +import ( + "context" + "encoding/json" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/accept" + "tidepool/internal/consume" + "tidepool/internal/materialize" + "tidepool/internal/store" +) + +// TASK 17d — THE OPT-OUT LIFECYCLE, SOFT TIER. +// +// Federation is DEFAULT-ON, so this record is the only way a user says stop — +// and what "stop" has to mean is two facts at once: nothing new goes out, AND +// nothing already queued goes out either. Today only the first is applied. The +// preference is recorded and the actor mirror is flipped, and the deliveries the +// user's own earlier posts left in the queue keep leaving, minutes or hours +// after they asked us not to. +// +// The two writes must land TOGETHER. Half of this state is a silent wrong +// answer in either direction: an actor disabled with live deliveries keeps +// publishing for someone who opted out, and cancelled deliveries under an actor +// still marked enabled loses the user's queued work while the bridge believes it +// is still federating for them. Nothing retries either half — the record is +// applied once, under a rev gate that rejects the replay. +// +// WHAT IS NOT REVOKED: the actor DOCUMENT. Peers that already hold this user's +// ids must still be able to resolve them, or every existing thread on the +// fediverse side breaks its author reference. Opting out stops the bridge +// SPEAKING for someone; it does not retract who they were. That is the 410 in +// the destructive tier, and the difference between the tiers is the whole +// design. + +const ( + // The opt-out record. rkey "self" is the lexicon's, and the rev is reused + // verbatim on the re-enable so the two commits are ordered by the gate the + // way a real client's would be. + odOptOutRev = "3lzodrev000001" + odReEnableRev = "3lzodrev000002" + + // Posts that leave queued work in two communities for one author, and one + // for the OTHER author — the axis that tells "cancel this actor's work" from + // "cancel this community's". + odPostInARKey = "3lzodpost00001" + odPostInBRKey = "3lzodpost00002" + odOtherRKey = "3lzodpost00003" + odAfterRKey = "3lzodpost00004" +) + +// TestOptingOutDisablesTheActorAndCancelsItsQueuedWork is the OUTER CONTRACT for +// 17d's soft tier. +// +// GIVEN a native author with pending deliveries to TWO communities, WHEN they +// write enabled=false, THEN the actor is disabled AND every pending delivery of +// theirs is cancelled, the other author's work is untouched, their actor +// document is still served, nothing is enqueued — and when they re-enable, +// admission resumes. +func TestOptingOutDisablesTheActorAndCancelsItsQueuedWork(t *testing.T) { + h := newHarness(t) + ctx := context.Background() + world := newModerationWorld(t, h) + + // --- GIVEN: queued work in both communities, for two different authors. + admitPost(t, world, mtAuthorDID, odPostInARKey, world.communityADID, "3lzodrev000010", 1_775_000_020_000_001) + admitPost(t, world, mtAuthorDID, odPostInBRKey, world.communityBDID, "3lzodrev000011", 1_775_000_020_000_002) + admitPost(t, world, mtCommenterDID, odOtherRKey, world.communityADID, "3lzodrev000012", 1_775_000_020_000_003) + + requireEveryDelivery(t, h.db, mtAuthorDID, groupID, "pending", + "precondition: the opting-out author has queued work for A") + requireEveryDelivery(t, h.db, mtAuthorDID, mtCommunityBAPID, "pending", + "precondition: and for B — one community could not tell 'this actor's work' from "+ + "'this community's'") + requireEveryDelivery(t, h.db, mtCommenterDID, groupID, "pending", + "precondition: and the OTHER author has queued work of their own") + + activitiesBefore := rowCount(t, h.db, "outbound_activities") + deliveriesBefore := rowCount(t, h.db, "outbound_deliveries") + + // --- WHEN: they opt out. Soft tier: no deleteRemote, nothing inferred. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, mtAuthorDID, odOptOutRev, "create", false, false, 1_775_000_021_000_001))) + + // --- THEN: the actor is disabled... + actor, err := store.NewAPActors(h.db).GetByDID(ctx, mtAuthorDID) + require.NoError(t, err, "the actor still EXISTS: opting out disables it, it does not erase it") + assert.False(t, actor.Enabled, + "and it is disabled: this mirror is what the consumer, the admission gate and the "+ + "delivery claim all read to answer 'may we still speak for this user'") + + // ...and every delivery they had queued is cancelled, in BOTH communities. + assertEveryDelivery(t, h.db, mtAuthorDID, groupID, "cancelled", + "their queued work for A is cancelled: a user who has asked us to stop and then "+ + "watches their posts keep arriving on Lemmy for the next hour has been told no "+ + "twice — once by us, once by the queue") + assertEveryDelivery(t, h.db, mtAuthorDID, mtCommunityBAPID, "cancelled", + "and their work for B with it: the opt-out is about the ACTOR, so a cancellation "+ + "scoped to one community leaves them federating everywhere else they ever posted") + + // ...while the other author is untouched. + assertEveryDelivery(t, h.db, mtCommenterDID, groupID, "pending", + "the OTHER author's work stands: cancelling by community — or by anything but this "+ + "actor — silences people who asked for nothing") + + // --- AND: their actor document is STILL SERVED. + doc, err := h.client.FetchActor(ctx, mtUserOrigin+"/ap/actor/"+mtAuthorDID) + require.NoError(t, err, + "the actor document must still resolve: every Note and Page this bridge already "+ + "delivered names this actor, and a 404 breaks the author reference on every one "+ + "of them — retracting the identity is the DESTRUCTIVE tier, and the difference "+ + "between the two is what the user chose") + assert.Equal(t, mtUserOrigin+"/ap/actor/"+mtAuthorDID, doc.ID) + + // --- AND: nothing went out about it. + assert.Equal(t, activitiesBefore, rowCount(t, h.db, "outbound_activities"), + "NO outbound activity: the soft tier tells peers nothing — a Delete{Person} here "+ + "would be the destructive tier applied to a user who did not ask for it, and no "+ + "peer un-deletes") + assert.Equal(t, deliveriesBefore, rowCount(t, h.db, "outbound_deliveries"), + "and no new delivery; the cancellation changes state, it does not add rows") + + // --- WHEN: they change their mind. + require.NoError(t, world.dispatcher.HandleEvent(ctx, + odFederationEvent(t, mtAuthorDID, odReEnableRev, "create", true, false, 1_775_000_022_000_001))) + + reenabled, err := store.NewAPActors(h.db).GetByDID(ctx, mtAuthorDID) + require.NoError(t, err) + assert.True(t, reenabled.Enabled, "the actor is enabled again, under its ORIGINAL identity") + + admitPost(t, world, mtAuthorDID, odAfterRKey, world.communityADID, "3lzodrev000013", 1_775_000_023_000_001) + afterURI := "at://" + mtAuthorDID + "/" + materialize.CollectionPostV2 + "/" + odAfterRKey + status, _ := admissionFor(t, h.db, world.communityADID, afterURI) + assert.Equal(t, accept.StatusAccepted, status, + "and admission resumes: re-enabling is the same user under the same actor, so a "+ + "soft opt-out has to be fully reversible — that is what makes it the tier a user "+ + "can safely choose") + _, _, err = h.manager.GetRecord(ctx, + world.communityADID, materialize.CollectionAcceptance, testDigestRKey(afterURI)) + assert.NoError(t, err, "with the acceptance to prove it (err=%v)", err) + + // The cancelled deliveries stay cancelled. Re-enabling restores the user's + // ability to federate; it does not resurrect work they cancelled by asking + // us to stop — those posts were withdrawn, and re-sending them would publish + // on their behalf something they had already taken back. + assertEveryDelivery(t, h.db, mtAuthorDID, mtCommunityBAPID, "cancelled", + "and the withdrawn work stays withdrawn: a re-enable is 'federate me from now on', "+ + "not 'replay what I stopped'") +} + +// odFederationEvent builds a social.coves.bridge.federation commit. +func odFederationEvent(t *testing.T, did, rev, operation string, enabled, deleteRemote bool, timeUS int64) *consume.JetstreamEvent { + t.Helper() + record := "" + if operation != "delete" { + record = fmt.Sprintf(`,"cid":%q,"record":{"$type":%q,"enabled":%t,"deleteRemote":%t}`, + mtPostCID, consume.CollectionFederation, enabled, deleteRemote) + } + frame := fmt.Sprintf( + `{"did":%q,"time_us":%d,"kind":"commit","commit":{"rev":%q,"operation":%q,`+ + `"collection":%q,"rkey":"self"%s}}`, + did, timeUS, rev, operation, consume.CollectionFederation, record) + var event consume.JetstreamEvent + require.NoError(t, json.Unmarshal([]byte(frame), &event), "the frame must be valid wire JSON") + return &event +} diff --git a/internal/optout/terminator.go b/internal/optout/terminator.go new file mode 100644 index 0000000..66f5d22 --- /dev/null +++ b/internal/optout/terminator.go @@ -0,0 +1,155 @@ +// Package optout implements the TERMINAL tier of the federation lifecycle +// (PLAN.md decision 19): what the bridge does when a native account is gone. +// +// It is a separate package from the consumer that calls it because the two make +// opposite mistakes. The consumer's job is to apply what a firehose event says; +// this tier's job is to DISBELIEVE it until the identity's own sources agree, +// because the action on the other side — asking every peer to delete a user's +// content — is one no peer undoes. +package optout + +import ( + "context" + "fmt" + "log/slog" + + "tidepool/internal/errors" + "tidepool/internal/store" +) + +// AccountConfirmer answers whether a DID's repo is really deleted, read from +// the identity's own sources rather than from the event that brought us here. +// *consume.HandleResolver implements it. +// +// THE THREE OUTCOMES ARE THE INTERFACE, and collapsing any two is the bug this +// seam exists to prevent: +// +// (true, nil) — confirmed deleted; the destructive tier may run. +// (false, nil) — confirmed live; nothing destructive, and the event is done. +// (_, err) — UNKNOWN; nothing destructive, nothing recorded, retry. +// +// "We could not confirm" is not "confirmed not deleted". That collapse is +// 17c-3's P1-c one layer up — a ban's expiry where absent and unparseable both +// became nil, nil meant permanent, and an author was excluded forever because a +// timestamp did not parse — and here it goes wrong in both directions at once: +// read as live it silently drops a real deletion and leaves the user's content +// federated forever, read as deleted it erases a user who never left. +type AccountConfirmer interface { + AccountStatus(ctx context.Context, did string) (deleted bool, err error) +} + +// RemoteContentDeleter asks peers to drop a user's already-federated content — +// the irreversible half. Optional: a deployment without it records the +// preference and leaves the work in the backlog. +type RemoteContentDeleter interface { + DeleteRemoteContent(ctx context.Context, did string) error +} + +// Options configures a Terminator. Confirmer and Prefs are required. +type Options struct { + // Confirmer re-verifies the deletion against PLC and the PDS. + Confirmer AccountConfirmer + // Prefs records the user's terminal state, so the intent survives a crash + // between deciding and sending. + Prefs store.FederationPrefs + // Deleter is the destructive seam. Optional (see RemoteContentDeleter). + Deleter RemoteContentDeleter + Logger *slog.Logger +} + +// Terminator applies decision 19's terminal tier: confirm, then act. +type Terminator struct { + confirmer AccountConfirmer + prefs store.FederationPrefs + deleter RemoteContentDeleter + logger *slog.Logger +} + +// NewTerminator validates options and builds a Terminator. +func NewTerminator(opts Options) (*Terminator, error) { + if opts.Confirmer == nil { + // REQUIRED. A terminator that cannot confirm would act on the event + // alone, which is the one thing this tier exists not to do. + return nil, errors.NewValidationError("confirmer", "must not be nil") + } + if opts.Prefs == nil { + return nil, errors.NewValidationError("prefs", "must not be nil") + } + logger := opts.Logger + if logger == nil { + logger = slog.Default() + } + if opts.Deleter == nil { + // Announced once, at construction, rather than discovered from a + // confirmed deletion quietly doing nothing. + logger.Warn("account terminator has no destructive seam wired: confirmed deletions " + + "will be recorded and left in the backlog") + } + return &Terminator{ + confirmer: opts.Confirmer, + prefs: opts.Prefs, + deleter: opts.Deleter, + logger: logger, + }, nil +} + +// TerminateAccount handles a repo whose #account status says deleted. +// +// CONFIRM FIRST, ALWAYS. The event is a claim about a moment that may have +// passed: reconnects replay events, a user can reactivate, and a deletion that +// was true an hour ago may not be true now. Nothing here is recoverable once +// sent, so the sequence is confirm → record → send, and each step only runs +// because the one before it succeeded. +// +// An unconfirmable status returns an ERROR, which is what leaves the consumer's +// per-DID seq un-advanced and the event redrivable. A confirmed LIVE account +// returns nil: there is nothing to do, and the seq must advance or a stale +// deletion event would wedge every later status change for that DID behind it. +func (t *Terminator) TerminateAccount(ctx context.Context, did string) error { + deleted, err := t.confirmer.AccountStatus(ctx, did) + if err != nil { + // UNKNOWN. Not "not deleted" — see AccountConfirmer. + return fmt.Errorf("confirm account deletion for %s: %w", did, err) + } + if !deleted { + // The event said deleted and the identity says otherwise, which is + // exactly why the confirm exists: a rewound cursor replaying a + // deletion the user has since reversed, or a status that never meant + // what the event implied. Logged at INFO because it is a real + // divergence between the firehose and PLC, and silence here would make + // 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 + } + + // RECORDED BEFORE ANYTHING IS SENT. Peers that honour a Delete cannot + // restore what they dropped, so the user's terminal state has to survive a + // crash between deciding and sending — and it is what stops the bridge + // federating for them again in the meantime. + if _, err := t.prefs.Upsert(ctx, store.FederationPref{ + DID: did, + Enabled: false, + DeleteRemote: true, + Source: store.FederationPrefSourceAccount, + }); err != nil { + return fmt.Errorf("record account deletion for %s: %w", did, err) + } + + if t.deleter == nil { + // The same shape the opt-out path already uses for an unwired + // destructive tier: the request is recorded, so it can be acted on from + // the backlog rather than the user's deletion being lost. Deliberately + // NOT degraded into a pause — a deletion half-handled as a pause looks + // handled in the database and is not. + t.logger.Warn("account confirmed deleted but no destructive seam is wired", + slog.String("did", did)) + return nil + } + if err := t.deleter.DeleteRemoteContent(ctx, did); err != nil { + return fmt.Errorf("delete remote content for %s: %w", did, err) + } + t.logger.Info("account confirmed deleted; remote content withdrawal requested", + slog.String("did", did)) + return nil +} diff --git a/internal/outbound/enqueuer.go b/internal/outbound/enqueuer.go index 5cde024..c70d8b5 100644 --- a/internal/outbound/enqueuer.go +++ b/internal/outbound/enqueuer.go @@ -114,24 +114,33 @@ func (e *Enqueuer) EnqueueActivity(ctx context.Context, tx *sql.Tx, actorDID, or } // The activity is the atomicity anchor: ON CONFLICT DO NOTHING, so a - // redelivered intent re-derives the same id and reports inserted=false. - // Because the activity, its mapping and its delivery all ride ONE - // transaction, an already-present activity id already has the other two — - // so a re-enqueue returns here without re-inserting the delivery (which - // would violate its PK and poison the tx). - inserted, err := e.activities.InsertTx(ctx, tx, store.OutboundActivity{ + // redelivered intent re-derives the same id and simply finds it there. The + // canonical payload is never rewritten — a peer may already hold it. + // + // IT DOES NOT DECIDE WHETHER THE DELIVERY EXISTS, and it used to: an early + // return here on an already-present activity assumed one activity meant one + // delivery, which is true of everything addressed to a single community and + // false of the one case the fan-out schema was built for. Delete{Person} is + // ONE activity to EVERY inbox an actor reached, so the second leg found the + // activity present and returned before resolving an inbox or writing + // anything — an erasure request reaching one instance out of many, with an + // activity row, a delivery row, a clean worker and clean metrics. It bit on + // REDELIVERY rather than first send, which is exactly when the destructive + // tier runs. + // + // The two questions are now asked separately: this one is "does the activity + // exist", the delivery's own idempotent insert below is "does THIS delivery + // exist" (EnqueueTx, which returns the standing row rather than violating + // its PK and poisoning this transaction). + if _, err := e.activities.InsertTx(ctx, tx, store.OutboundActivity{ ActivityID: intent.ActivityID(), ActorDID: actorDID, Kind: translated.Kind, Payload: translated.Payload, ParentATURI: parentATURI, - }) - if err != nil { + }); err != nil { return fmt.Errorf("insert outbound activity %s: %w", intent.ActivityID(), err) } - if !inserted { - return nil - } // The object mapping (bridge-origin) makes GET /ap/object serve the record. // Votes have no servable object and a self-delete maps nothing new. @@ -160,6 +169,63 @@ func (e *Enqueuer) EnqueueActivity(ctx context.Context, tx *sql.Tx, actorDID, or return nil } +// EnqueueFanOut writes ONE activity and a delivery for EVERY target — the shape +// EnqueueActivity cannot express, because that one derives its single inbox from +// the intent's community and this one is addressed to instances rather than to a +// community. +// +// It is the destructive tier's entry point: Delete{Person} goes to every inbox +// the actor's content reached, from the delivery history (store's +// DistinctInboxesForActor). The canonical payload is shared by all of them, so +// the peer sees one erasure request however many times it is addressed. +// +// Each target keeps its own ordering key so the withdrawal serializes on the +// line that instance's other traffic already uses. Duplicates are no-ops: the +// delivery insert returns the standing row rather than violating its PK, which +// is what makes a retried fan-out reach the inboxes it missed without disturbing +// the ones it did not. +// +// 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, + intent consume.Intent, targets []store.DeliveryTarget) error { + + if tx == nil { + return errors.NewValidationError("tx", "must not be nil") + } + if len(targets) == 0 { + return nil + } + actor, err := e.actors.GetByDID(ctx, actorDID) + if err != nil { + return fmt.Errorf("resolve actor %s: %w", actorDID, err) + } + translated, err := e.translator.Translate(actor.ActorID, intent) + if err != nil { + 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, + }); err != nil { + return fmt.Errorf("insert outbound activity %s: %w", intent.ActivityID(), err) + } + for _, target := range targets { + if _, err := e.deliveries.EnqueueTx(ctx, tx, store.OutboundDelivery{ + ActivityID: intent.ActivityID(), + TargetInbox: target.Inbox, + OrderingKey: target.OrderingKey, + }); err != nil { + return fmt.Errorf("enqueue delivery for %s to %s: %w", + intent.ActivityID(), target.Inbox, err) + } + } + return nil +} + // 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 new file mode 100644 index 0000000..698bffe --- /dev/null +++ b/internal/outbound/fanout_test.go @@ -0,0 +1,182 @@ +package outbound + +import ( + "context" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/consume" + "tidepool/internal/store" +) + +// TASK 17d — ONE ACTIVITY, MANY INBOXES. +// +// Every activity this bridge has sent so far goes to exactly one community, so +// the enqueuer's idempotency shortcut has been correct by coincidence: when the +// activity row is already there, its delivery is already there too, and +// returning early saves a PK violation. +// +// inserted, err := e.activities.InsertTx(...) +// if !inserted { return nil } // ← skips the delivery insert +// +// The destructive opt-out tier breaks that premise. `Delete{Person, +// removeData:true}` is ONE activity addressed to EVERY community this actor +// ever delivered to — the fan-out the outbound schema was built for (migration +// 020's comment names this exact case). The moment the second inbox is enqueued +// under an activity id the first one already inserted, the shortcut returns +// before writing anything, and the user's erasure request reaches ONE instance +// out of however many hold their content. +// +// The failure is silent and looks like success: an activity row exists, a +// delivery row exists, the worker delivers it, the metrics are clean. Only a +// count of delivery rows against the actor's delivery history tells the truth — +// and nothing counts that. For a request that cannot be re-asked (peers do not +// un-delete, and asking twice risks deleting content a re-enabled user has since +// restored) reaching one instance is worse than failing loudly. +// +// TWO PASSES, because they fail differently. The FIRST pass is the fan-out +// itself. The SECOND is the redelivery — and the redelivery is the one that +// matters most in production, since the destructive tier is precisely the work +// most likely to be retried after a crash. A test that enqueues once and counts +// rows would pass on an implementation that reaches one instance on every +// retry. +const ( + fanCommunityB = "https://lemmy.zip/c/technology" + fanInboxA = "https://lemmy.world/inbox" + fanInboxB = "https://lemmy.zip/inbox" +) + +// perCommunityInboxes answers a DIFFERENT inbox per community — the real shape, +// since each instance hosts its own. (Co-hosted communities share one, which is +// why the delivery PK is (activity, inbox) and the ordering key is the +// community: the two are independent axes.) +type perCommunityInboxes struct { + inboxes map[string]string + calledWith []string +} + +func (r *perCommunityInboxes) ResolveInbox(_ context.Context, communityAPID string) (string, error) { + r.calledWith = append(r.calledWith, communityAPID) + return r.inboxes[communityAPID], nil +} + +// TestEnqueuer_OneActivityReachesEveryInboxNotJustTheFirst is the fan-out +// contract the destructive tier rests on. +func TestEnqueuer_OneActivityReachesEveryInboxNotJustTheFirst(t *testing.T) { + conn := enqueuerTestDB(t) + ctx := context.Background() + seedEnqueuerActor(t, conn) + + resolver := &perCommunityInboxes{inboxes: map[string]string{ + outCommunityAPID: fanInboxA, + fanCommunityB: fanInboxB, + }} + enq := newTestEnqueuer(t, conn, resolver) + + // ONE activity id, addressed to two communities on two instances. The intent + // type is incidental — Delete{Person} has no intent yet, and the property is + // about the ACTIVITY and DELIVERY rows, which are shared by every intent. + base := commentEnqueueIntent(t) + fanOut := func(t *testing.T, pass string) { + t.Helper() + for _, community := range []string{outCommunityAPID, fanCommunityB} { + intent := base + intent.CommunityAPID = community + tx, err := conn.BeginTx(ctx, nil) + require.NoError(t, err) + require.NoError(t, + enq.EnqueueActivity(ctx, tx, outCommenterDID, outCommenterDID, "", intent), + "%s: enqueueing %s must not fail — an erasure request that errors half way "+ + "through its fan-out leaves the user's intent partly applied with no record "+ + "of which peers were reached", pass, community) + require.NoError(t, tx.Commit()) + } + } + + // --- PASS ONE: the fan-out. + fanOut(t, "first pass") + + assert.Equal(t, 1, count(t, conn, "outbound_activities"), + "ONE canonical activity: the same erasure request, not one per peer — the id is what "+ + "makes a redelivery recognisable to the peer as the same activity") + assert.Equal(t, 2, count(t, conn, "outbound_deliveries"), + "and a delivery PER INBOX: the fan-out schema exists for exactly this, and one row "+ + "for two instances means the user's content stays live on every instance but the "+ + "first — silently, with every metric reading clean") + + require.Contains(t, resolver.calledWith, fanCommunityB, + "the second community's inbox must even be RESOLVED: the early return fires before "+ + "the resolver is consulted, so a fan-out that never asks is one that never intended "+ + "to deliver") + + for _, target := range []struct{ inbox, orderingKey string }{ + {fanInboxA, outCommunityAPID}, + {fanInboxB, fanCommunityB}, + } { + delivery, err := store.NewOutboundDeliveries(conn).Get(ctx, base.ID, target.inbox) + require.NoError(t, err, + "a delivery must exist for %s: this is the row that carries the request to that "+ + "instance, and its absence is the erasure silently not happening there", + target.inbox) + assert.Equal(t, store.DeliveryStatePending, delivery.State) + assert.Equal(t, target.orderingKey, delivery.OrderingKey, + "addressed on that community's own serial line: the inbox is shared by co-hosted "+ + "communities, so the ordering key is what keeps the lines independent") + } + + // --- PASS TWO: the redelivery. Same activity id, same inboxes, again. + // + // This is the shape production actually meets — the destructive tier is the + // work most likely to be retried — and it is where the idempotency the early + // return exists to provide has to hold WITHOUT costing the second inbox. + fanOut(t, "redelivery") + + assert.Equal(t, 1, count(t, conn, "outbound_activities"), + "still ONE activity after a redelivery: the id is deterministic, and a second row "+ + "would present the peer with a second erasure to apply") + assert.Equal(t, 2, count(t, conn, "outbound_deliveries"), + "and still exactly TWO deliveries — no duplicates, none lost. Idempotency and "+ + "completeness are not in tension here: the delivery PK is (activity, inbox), so "+ + "re-enqueueing the same pair is a no-op while a NEW pair is a row that must be "+ + "written") +} + +// TestEnqueuer_ARedeliveredActivityToTheSameInboxStaysOneRow keeps the reason +// the early return exists. +// +// The shortcut is not arbitrary: the delivery PK is (activity_id, target_inbox), +// so re-inserting the same pair inside the caller's transaction would raise a +// unique violation and poison it — taking down the rev-gate advance riding the +// same tx with it. Whatever replaces the shortcut has to keep this true, or the +// fix for a silent under-delivery becomes a loud failure on every ordinary +// redelivery. +func TestEnqueuer_ARedeliveredActivityToTheSameInboxStaysOneRow(t *testing.T) { + conn := enqueuerTestDB(t) + ctx := context.Background() + seedEnqueuerActor(t, conn) + + enq := newTestEnqueuer(t, conn, &fakeInboxResolver{inbox: outSharedInbox}) + intent := commentEnqueueIntent(t) + + for pass := 1; pass <= 2; pass++ { + tx, err := conn.BeginTx(ctx, nil) + require.NoError(t, err) + require.NoError(t, enq.EnqueueActivity(ctx, tx, outCommenterDID, outCommunityAPID, outRootATURI, intent), + "pass %d must not error: a redelivered intent re-derives the same activity id, and "+ + "a unique violation here would poison the caller's transaction — the rev-gate "+ + "advance rides it, so the event would replay forever", pass) + require.NoError(t, tx.Commit()) + } + + assert.Equal(t, 1, count(t, conn, "outbound_activities"), "one activity") + assert.Equal(t, 1, count(t, conn, "outbound_deliveries"), + "and ONE delivery: the same (activity, inbox) pair twice is the same delivery, and a "+ + "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{} diff --git a/internal/outbound/purge.go b/internal/outbound/purge.go new file mode 100644 index 0000000..2c55dbe --- /dev/null +++ b/internal/outbound/purge.go @@ -0,0 +1,193 @@ +package outbound + +import ( + "context" + "database/sql" + "fmt" + "log/slog" + + "tidepool/internal/consume" + "tidepool/internal/errors" + "tidepool/internal/store" +) + +// Purger is task 17d's DESTRUCTIVE opt-out tier (decision 11's second tier): the +// user asked not merely that the bridge stop speaking for them, but that what it +// already said be withdrawn. It satisfies consume.RemoteContentDeleter. +// +// THREE CONSEQUENCES, AND EACH IS UNREACHABLE FROM THE OTHERS: +// +// 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. +// 3. The actor document stops resolving — 410 Gone. +// +// IRREVERSIBLE, and reached only from an explicit enabled=false + +// deleteRemote=true record or a CONFIRMED account deletion — never inferred. +// Lemmy un-deletes a person on refetch; other software does not, so nothing here +// promises resurrection and no code path clears the tombstone. +// +// ONE TRANSACTION. A withdrawal that half-applied is the worst of both tiers: +// an actor tombstoned with no Delete sent is an identity that vanished while its +// content stayed, and a Delete sent with the votes left standing is a purge that +// keeps voting. Everything below either commits together or not at all — and +// because it is reached through the consumer, a rollback leaves the record's rev +// gate un-advanced and the whole request replayable. +type Purger struct { + db *sql.DB + userOrigin string + enqueuer *Enqueuer + deliveries store.OutboundDeliveries + votes store.OutboundVotes + actors store.APActors + communities store.Communities + logger *slog.Logger +} + +// NewPurger builds the destructive tier. The stores are defaulted from db, as +// every other constructor in this package does. +func NewPurger(db *sql.DB, userOrigin string, enqueuer *Enqueuer) *Purger { + return &Purger{ + db: db, + userOrigin: userOrigin, + enqueuer: enqueuer, + deliveries: store.NewOutboundDeliveries(db), + votes: store.NewOutboundVotes(db), + actors: store.NewAPActors(db), + communities: store.NewCommunities(db), + logger: slog.Default(), + } +} + +// WithLogger returns the purger with a logger attached. Optional: the default is +// slog.Default(), and every line this tier writes is worth having. +func (p *Purger) WithLogger(logger *slog.Logger) *Purger { + if logger != nil { + p.logger = logger + } + return p +} + +// personDeleteOp is the activity-id op string for a person withdrawal. It is +// distinct from every record operation so the deterministic id can never collide +// with an activity about the user's CONTENT — the preimage is +// (tag, subject, op, seq), and here the subject is the actor rather than a +// record. +const personDeleteOp = "delete-person" + +// DeleteRemoteContent asks every peer that ever received this actor's content to +// delete it, retracts the votes they still hold, and stops serving the actor. +// +// IDEMPOTENT by construction rather than by a guard: the activity id is +// deterministic, the delivery insert returns the standing row instead of +// duplicating it, and the tombstone keeps its original timestamp. A retried +// purge therefore reaches inboxes the first attempt missed without re-sending +// anything to the ones it reached — which matters because this is the work most +// likely to be retried after a crash, and asking a peer twice risks deleting +// content a re-enabled user has since restored. +func (p *Purger) DeleteRemoteContent(ctx context.Context, did string) error { + if did == "" { + 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. + targets, err := p.deliveries.DistinctInboxesForActor(ctx, did) + if err != nil { + return err + } + liveVotes, err := p.votes.ListDeliveredForActor(ctx, did) + if err != nil { + return err + } + + tx, err := p.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("purge %s: begin: %w", did, err) + } + defer func() { _ = tx.Rollback() }() + + 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 + // activity it already has. + ID: consume.ActivityID(p.userOrigin, did, personDeleteOp, 0), + }, targets); err != nil { + return err + } + + if err := p.undoLiveVotes(ctx, tx, did, liveVotes); err != nil { + return err + } + + // 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 + // it announces has been written. + if err := p.actors.TombstoneTx(ctx, tx, did); err != nil && !errors.IsNotFound(err) { + return fmt.Errorf("tombstone actor %s: %w", did, err) + } + + if err := tx.Commit(); err != nil { + return fmt.Errorf("purge %s: commit: %w", did, err) + } + 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))) + return nil +} + +// undoLiveVotes retracts the votes peers still hold and marks them retracted. +// +// The state flip is the half that is easy to miss and impossible to see: without +// it the reseed keeps subtracting these votes from the origin's API tally for an +// actor that no longer exists, so a subject's served score drifts down and stays +// there. It is written HERE rather than waiting for the Undo's delivery +// callback, deliberately: a purge is terminal, and a withdrawn actor's vote must +// stop counting when we decide to withdraw it, not if and when a peer confirms. +// 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) + } + + // 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 + // is preserved by the upsert, which matters: it is the id the Like went out + // under and the Undo has to embed it. + vote.DeliveredState = store.DeliveredStateUndone + bumped, err := p.votes.UpsertTx(ctx, tx, vote) + if err != nil { + return fmt.Errorf("retract vote %s: %w", vote.VoteATURI, err) + } + + if err := p.enqueuer.EnqueueActivity(ctx, tx, did, did, vote.SubjectATURI, consume.VoteIntent{ + Op: consume.OperationUndo, + VoteATURI: vote.VoteATURI, + SubjectAPID: vote.SubjectAPID, + // Read back from state, never guessed: an Undo{Like} withdrawing a + // Dislike would move the peer's count the wrong way. + Direction: vote.Direction, + ID: consume.ActivityID(p.userOrigin, vote.VoteATURI, consume.OperationUndo, bumped.ActivitySeq), + InnerActivityID: vote.CurrentActivityID, + CommunityAPID: community.APGroupID, + }); err != nil { + return err + } + } + return nil +} diff --git a/internal/outbound/translator.go b/internal/outbound/translator.go index 6c24108..c9966f9 100644 --- a/internal/outbound/translator.go +++ b/internal/outbound/translator.go @@ -64,6 +64,8 @@ func (t *Translator) Translate(actorID string, intent consume.Intent) (*Translat return t.post(actorID, typed) case consume.VoteIntent: return t.vote(actorID, typed) + case consume.PersonDeleteIntent: + return t.personDelete(actorID, typed) default: return nil, errors.NewValidationError("intent", fmt.Sprintf("no translation for %T", intent)) } @@ -203,6 +205,36 @@ func (t *Translator) wrapActivity(kind, id, actorID, communityAPID string, objec } } +// personDelete renders the destructive opt-out: Delete{Person, removeData:true}. +// +// The ACTOR AND THE OBJECT ARE THE SAME ID, because this is the user withdrawing +// themselves — the one activity this bridge sends that is about its own sender. +// +// removeData is what makes it a purge rather than a marker. Verified against +// Lemmy 0.19: without it, Delete{Person} marks the person deleted and KEEPS +// their content; with it, the content is purged. A destructive tier that omitted +// it would leave every post standing while reporting success. +// +// It carries NO cc and NO audience, unlike every other activity here. Those name +// the one community an activity is for, and this one is for every instance the +// user ever reached — one canonical payload, fanned out over many inboxes, so it +// cannot name any of them. `to` is Public, which is what the addressing means +// when the audience is "everyone who holds this identity". +func (t *Translator) personDelete(actorID string, intent consume.PersonDeleteIntent) (*TranslatedActivity, error) { + if intent.ActorDID == "" { + return nil, errors.NewValidationError("actorDid", "must not be empty") + } + return t.finish(intent.ID, "Delete", "", map[string]any{ + "@context": contextActivityStreams, + "id": intent.ID, + "type": "Delete", + "actor": actorID, + "object": actorID, + "removeData": true, + "to": []string{ap.PublicAudience}, + }) +} + // deleteActivity is a self-delete: a Delete of the bare object URL and NOTHING // else. Lemmy reads a `summary` on a Delete as a MOD-REMOVAL reason, so a // self-delete must omit it or it looks like moderation. diff --git a/internal/personas/serving.go b/internal/personas/serving.go index 61cb1cc..62cdeb0 100644 --- a/internal/personas/serving.go +++ b/internal/personas/serving.go @@ -223,6 +223,21 @@ func (s *Service) handleActorDocument(w http.ResponseWriter, r *http.Request, di http.NotFound(w, r) return } + // WITHDRAWN (task 17d's destructive tier): 410 Gone, and specifically not + // 404. Gone is the statement "this existed and was withdrawn", which is what + // lets a peer stop re-fetching and clean up its own copy; 404 reads as + // "never heard of them", which several implementations treat as a transient + // lookup failure and retry indefinitely. + // + // This is the ONE observable difference between the two opt-out tiers, which + // is why it is a status code rather than a flag: a soft-disabled actor's + // document keeps resolving, because every Note and Page already delivered + // names it and revoking it would orphan the author reference on every + // existing thread. + if actor.IsTombstoned() { + http.Error(w, "gone", http.StatusGone) + return + } origin, err := actorOrigin(actor) if err != nil { diff --git a/internal/store/ap_actors.go b/internal/store/ap_actors.go index 9e73246..1cb3207 100644 --- a/internal/store/ap_actors.go +++ b/internal/store/ap_actors.go @@ -44,6 +44,12 @@ type APActor struct { // DeliveryPaused is the transient #account state (decision 19): // delivery stops, identity stays. DeliveryPaused bool + // TombstonedAt is set when the DESTRUCTIVE tier withdrew this identity + // (task 17d): the actor document answers 410 Gone from then on. It is + // TERMINAL — no path clears it — because peers that honoured the Delete + // cannot restore what they dropped, so nothing here may promise + // resurrection. + TombstonedAt *time.Time // Profile cache, refreshed from the appview/PDS (task 14 owns sync). DisplayName string Summary string @@ -72,7 +78,7 @@ func NewAPActors(db *sql.DB) APActors { const apActorColumns = ` did, kind, actor_id, normalized_origin, local_part, rsa_key_sealed, rsa_key_version, public_key_pem, - enabled, enabled_at, disabled_at, delivery_paused, + enabled, enabled_at, disabled_at, delivery_paused, tombstoned_at, display_name, summary, avatar_url, created_at, updated_at` func (r *postgresAPActors) Create(ctx context.Context, actor APActor) (*APActor, error) { @@ -144,6 +150,17 @@ func (r *postgresAPActors) GetByOriginLocalPart(ctx context.Context, normalizedO } func (r *postgresAPActors) SetEnabled(ctx context.Context, did string, enabled bool) error { + return r.setEnabled(ctx, r.db, did, enabled) +} + +func (r *postgresAPActors) SetEnabledTx(ctx context.Context, tx *sql.Tx, did string, enabled bool) error { + if tx == nil { + return errors.NewValidationError("tx", "must not be nil") + } + return r.setEnabled(ctx, tx, did, enabled) +} + +func (r *postgresAPActors) setEnabled(ctx context.Context, ex execer, did string, enabled bool) error { // Both transitions are stamped, and disabled_at is CLEARED on re-enable // so the column answers "is this actor currently disabled, and since // when" rather than "was it ever disabled". enabled_at is not cleared on @@ -155,7 +172,41 @@ func (r *postgresAPActors) SetEnabled(ctx context.Context, did string, enabled b disabled_at = CASE WHEN $2 THEN NULL ELSE now() END, updated_at = now() WHERE did = $1` - return r.execOne(ctx, "set enabled", did, query, did, enabled) + return execOneRow(ctx, ex, "set enabled", did, query, did, enabled) +} + +// IsTombstoned reports whether the destructive tier withdrew this identity. +func (a *APActor) IsTombstoned() bool { return a.TombstonedAt != nil } + +func (r *postgresAPActors) Tombstone(ctx context.Context, did string) error { + return r.tombstone(ctx, r.db, did) +} + +func (r *postgresAPActors) TombstoneTx(ctx context.Context, tx *sql.Tx, did string) error { + if tx == nil { + return errors.NewValidationError("tx", "must not be nil") + } + return r.tombstone(ctx, tx, did) +} + +func (r *postgresAPActors) tombstone(ctx context.Context, ex execer, did string) error { + // COALESCE keeps the ORIGINAL withdrawal time on a re-run, which makes the + // whole thing one idempotent statement: affected == 0 can only mean the + // actor does not exist. A destructive tier that re-stamped would lose the + // only record of when the user was actually withdrawn. + // + // enabled is set false in the same statement, because a tombstoned actor + // that still resolved through webfinger would be discoverable after being + // withdrawn — the two say the same thing from different sides and must not + // be able to disagree. + query := ` + UPDATE ap_actors SET + tombstoned_at = COALESCE(tombstoned_at, now()), + enabled = FALSE, + disabled_at = COALESCE(disabled_at, now()), + updated_at = now() + WHERE did = $1` + return execOneRow(ctx, ex, "tombstone", did, query, did) } func (r *postgresAPActors) SetPaused(ctx context.Context, did string, paused bool) error { @@ -180,7 +231,13 @@ func (r *postgresAPActors) UpdateProfile(ctx context.Context, did string, profil // execOne runs a single-row mutator and reports a missed DID as NotFound. func (r *postgresAPActors) execOne(ctx context.Context, what, did, query string, args ...any) error { - result, err := r.db.ExecContext(ctx, query, args...) + return execOneRow(ctx, r.db, what, did, query, args...) +} + +// execOneRow runs a single-row UPDATE on either the pool or a caller's +// transaction, mapping "matched nothing" to NotFound. +func execOneRow(ctx context.Context, ex execer, what, did, query string, args ...any) error { + result, err := ex.ExecContext(ctx, query, args...) if err != nil { return fmt.Errorf("%s for ap_actor %q: %w", what, did, err) } @@ -197,12 +254,12 @@ func (r *postgresAPActors) execOne(ctx context.Context, what, did, query string, func scanAPActor(row rowScanner) (*APActor, error) { var actor APActor var kind string - var enabledAt, disabledAt sql.NullTime + var enabledAt, disabledAt, tombstonedAt sql.NullTime err := row.Scan( &actor.DID, &kind, &actor.ActorID, &actor.NormalizedOrigin, &actor.LocalPart, &actor.RSAKeySealed, &actor.RSAKeyVersion, &actor.PublicKeyPEM, - &actor.Enabled, &enabledAt, &disabledAt, &actor.DeliveryPaused, + &actor.Enabled, &enabledAt, &disabledAt, &actor.DeliveryPaused, &tombstonedAt, &actor.DisplayName, &actor.Summary, &actor.AvatarURL, &actor.CreatedAt, &actor.UpdatedAt, ) @@ -216,5 +273,8 @@ func scanAPActor(row rowScanner) (*APActor, error) { if disabledAt.Valid { actor.DisabledAt = &disabledAt.Time } + if tombstonedAt.Valid { + actor.TombstonedAt = &tombstonedAt.Time + } return &actor, nil } diff --git a/internal/store/interfaces.go b/internal/store/interfaces.go index ac5b9f6..d7d85a1 100644 --- a/internal/store/interfaces.go +++ b/internal/store/interfaces.go @@ -215,10 +215,36 @@ type APActors interface { // missing actor is an error satisfying errors.IsNotFound. SetEnabled(ctx context.Context, did string, enabled bool) error + // SetEnabledTx is SetEnabled on an existing transaction — the seam the + // opt-out uses so the flag and the cancellation of the actor's queued work + // land together. Half of that state is a silent wrong answer in either + // direction: an actor disabled with live deliveries keeps publishing for + // someone who opted out, and cancelled deliveries under an enabled actor + // lose work while the bridge believes it still federates for them. A nil tx + // is an error satisfying errors.IsValidation. + SetEnabledTx(ctx context.Context, tx *sql.Tx, did string, enabled bool) error + // SetPaused toggles delivery_paused (the transient #account state). // A missing actor is an error satisfying errors.IsNotFound. SetPaused(ctx context.Context, did string, paused bool) error + // Tombstone withdraws the identity (task 17d's destructive tier): the actor + // document answers 410 Gone from then on, and the actor is disabled in the + // same statement so it cannot stay discoverable after being withdrawn. + // + // TERMINAL. Nothing clears it, and nothing here promises resurrection: + // peers that honoured the Delete this accompanies cannot restore what they + // dropped. Re-tombstoning preserves the original time. A missing actor is an + // error satisfying errors.IsNotFound. + Tombstone(ctx context.Context, did string) error + + // TombstoneTx is Tombstone on an existing transaction — the seam the + // destructive tier uses so the withdrawal and the activities that announce + // it commit together. An actor tombstoned without the Delete being enqueued + // is an identity that vanished while its content stayed. A nil tx is an + // error satisfying errors.IsValidation. + TombstoneTx(ctx context.Context, tx *sql.Tx, did string) error + // UpdateProfile refreshes the cached display name, summary, and avatar // and bumps updated_at. It NEVER touches local_part — the identity // handler (task 14) reaches this method on handle changes, and the @@ -490,6 +516,12 @@ 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) + // SetDeliveredState transitions the delivery state. An unknown state is // an error satisfying errors.IsValidation; a missing vote is an error // satisfying errors.IsNotFound. @@ -599,7 +631,18 @@ type OutboundDeliveries interface { Enqueue(ctx context.Context, delivery OutboundDelivery) (*OutboundDelivery, error) // EnqueueTx is Enqueue on an existing transaction — rides the enqueuer's - // gate tx. A nil tx is an error satisfying errors.IsValidation. + // gate tx — and is IDEMPOTENT where Enqueue refuses: a duplicate (activity, + // inbox) returns the STANDING row instead of an error. + // + // Both halves of that are load-bearing. A unique violation inside a caller's + // transaction aborts the whole transaction, so the rev-gate advance riding it + // dies too and the event replays forever. And a duplicate is not a caller + // bug here: ONE activity fans out to MANY inboxes, so a redelivery re-visits + // pairs that already exist while others still need writing. Returning the + // standing row (never resetting it) is what keeps a delivered, cancelled or + // poisoned delivery from being revived by a replay. + // + // A nil tx is an error satisfying errors.IsValidation. EnqueueTx(ctx context.Context, tx *sql.Tx, delivery OutboundDelivery) (*OutboundDelivery, error) // ClaimNext atomically claims the oldest processable delivery and @@ -639,8 +682,18 @@ type OutboundDeliveries interface { // cancelled (the consent/kill-switch withdrawal — a disabled or paused // actor's create/update work is parked, never poisoned). Terminal // deliveries are untouched. Returns how many rows were cancelled. + // + // It is scoped to the ACTOR across every community they have work in, + // because that is the scope of the decision: an opt-out cancelled per + // community would leave the user federating everywhere else they ever + // 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 + // an error satisfying errors.IsValidation. + CancelForActorTx(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. @@ -650,6 +703,14 @@ type OutboundDeliveries interface { // error satisfying errors.IsNotFound. Get(ctx context.Context, activityID, targetInbox string) (*OutboundDelivery, error) + // DistinctInboxesForActor lists every inbox this actor's content has been + // delivered to, one row per inbox, each carrying an ordering key from the + // history. It is the address book the destructive tier fans a Delete{Person} + // out over — the delivery history is the only record of which instances hold + // a user's content — and it is deliberately blind to delivery STATE (see the + // implementation for why, and for the instances it cannot reach). + DistinctInboxesForActor(ctx context.Context, actorDID string) ([]DeliveryTarget, error) + // ParentDeliveryPoisoned reports whether the delivery of the child's ACTUAL // parent (the activity that federated parentATURI as its object, to the same // inbox) is poisoned — the causal signal task 15's worker reads to poison a diff --git a/internal/store/models.go b/internal/store/models.go index 78c3faa..7d0a991 100644 --- a/internal/store/models.go +++ b/internal/store/models.go @@ -196,6 +196,16 @@ type CommunityBan struct { RemoveData bool } +// DeliveryTarget is one place an actor's content has already been delivered: +// the inbox that received it, and an ordering key that inbox's traffic is +// already serialized on. It is what a fan-out addresses — one activity, many +// targets — and it comes from the delivery history rather than from the +// communities table, because the question is where the content WENT. +type DeliveryTarget struct { + Inbox string + OrderingKey string +} + // ServiceKey is one of the bridge's own long-lived keys, keyed by purpose // name. KeyMaterial's encoding is per-row: plaintext PKCS#8 PEM for // "service-actor" (the AP-side RSA signing key — the bridge's own service @@ -251,12 +261,18 @@ const ( FederationPrefSourceRecord FederationPrefSource = "record" // FederationPrefSourceProbe means the bridge fetched the record itself. FederationPrefSourceProbe FederationPrefSource = "probe" + // FederationPrefSourceAccount means the preference was not expressed by the + // user at all: their ACCOUNT is gone, confirmed against PLC and the PDS by + // the terminal tier (decision 19). It is the one value that distinguishes a + // user who opted out — a decision they can reverse — from one who was + // deleted, which they cannot. + FederationPrefSourceAccount FederationPrefSource = "account" ) // Valid reports whether the value is a known source. func (s FederationPrefSource) Valid() bool { switch s { - case FederationPrefSourceRecord, FederationPrefSourceProbe: + case FederationPrefSourceRecord, FederationPrefSourceProbe, FederationPrefSourceAccount: return true } return false diff --git a/internal/store/outbound_deliveries.go b/internal/store/outbound_deliveries.go index 497a3a6..6aa8b1a 100644 --- a/internal/store/outbound_deliveries.go +++ b/internal/store/outbound_deliveries.go @@ -28,17 +28,6 @@ const deliveryColumns = ` last_error_class, response_excerpt, created_at, updated_at` func (r *postgresOutboundDeliveries) Enqueue(ctx context.Context, delivery OutboundDelivery) (*OutboundDelivery, error) { - return r.enqueue(ctx, r.db, delivery) -} - -func (r *postgresOutboundDeliveries) EnqueueTx(ctx context.Context, tx *sql.Tx, delivery OutboundDelivery) (*OutboundDelivery, error) { - if tx == nil { - return nil, errors.NewValidationError("tx", "must not be nil") - } - return r.enqueue(ctx, tx, delivery) -} - -func (r *postgresOutboundDeliveries) enqueue(ctx context.Context, q execer, delivery OutboundDelivery) (*OutboundDelivery, error) { // A fresh delivery is pending, unattempted, unclaimed: state, attempts, // next_attempt_at and seq are all defaulted by the table. A duplicate // (activity, inbox) pair violates the PK — mapped to AlreadyExists rather @@ -48,7 +37,7 @@ func (r *postgresOutboundDeliveries) enqueue(ctx context.Context, q execer, deli VALUES ($1, $2, $3) RETURNING` + deliveryColumns - stored, err := scanOutboundDelivery(q.QueryRowContext(ctx, query, + stored, err := scanOutboundDelivery(r.db.QueryRowContext(ctx, query, delivery.ActivityID, delivery.TargetInbox, delivery.OrderingKey)) if err != nil { if _, ok := uniqueViolation(err); ok { @@ -61,6 +50,56 @@ func (r *postgresOutboundDeliveries) enqueue(ctx context.Context, q execer, deli return stored, nil } +func (r *postgresOutboundDeliveries) EnqueueTx(ctx context.Context, tx *sql.Tx, delivery OutboundDelivery) (*OutboundDelivery, error) { + if tx == nil { + return nil, errors.NewValidationError("tx", "must not be nil") + } + // IDEMPOTENT, unlike its pool-backed sibling above, and the difference is + // deliberate on both sides. + // + // This one rides a CALLER'S transaction, where a unique violation is not an + // error the caller can inspect and move past — postgres aborts the whole + // transaction, taking the rev-gate advance riding it down too, so the event + // replays forever. And a duplicate here is not a caller bug: ONE activity + // fans out to MANY inboxes (Delete{Person} to every community an actor + // delivered to), so a redelivery legitimately re-enqueues pairs that already + // exist while others still need writing. + // + // The STANDING row wins. Re-enqueueing must never reset a delivery that has + // since been delivered, cancelled by an opt-out, or poisoned — the row's + // state is the record of what happened to it. That also settles the + // co-hosted case: two communities sharing one inbox collapse to a single + // delivery, keeping the ordering key of the first, because one POST to that + // inbox is one POST however many communities it serves. + query := ` + INSERT INTO outbound_deliveries (activity_id, target_inbox, ordering_key) + VALUES ($1, $2, $3) + ON CONFLICT (activity_id, target_inbox) DO NOTHING + RETURNING` + deliveryColumns + + stored, err := scanOutboundDelivery(tx.QueryRowContext(ctx, query, + delivery.ActivityID, delivery.TargetInbox, delivery.OrderingKey)) + if stderrors.Is(err, sql.ErrNoRows) { + // DO NOTHING returns no row, so the delivery was already there. Read it + // back on the same transaction: callers get the same contract either way + // — a row that exists — and never have to tell the two apart. + existing, rerr := scanOutboundDelivery(tx.QueryRowContext(ctx, + `SELECT`+deliveryColumns+` FROM outbound_deliveries + WHERE activity_id = $1 AND target_inbox = $2`, + delivery.ActivityID, delivery.TargetInbox)) + if rerr != nil { + return nil, fmt.Errorf("read existing outbound_delivery %q -> %q: %w", + delivery.ActivityID, delivery.TargetInbox, rerr) + } + return existing, nil + } + if err != nil { + return nil, fmt.Errorf("enqueue outbound_delivery %q -> %q: %w", + delivery.ActivityID, delivery.TargetInbox, err) + } + return stored, nil +} + func (r *postgresOutboundDeliveries) ClaimNext(ctx context.Context, lease time.Duration) (*OutboundDelivery, error) { if lease <= 0 { return nil, errors.NewValidationError("lease", "must be positive") @@ -225,9 +264,22 @@ func (r *postgresOutboundDeliveries) markResult(ctx context.Context, op, query s } func (r *postgresOutboundDeliveries) CancelForActor(ctx context.Context, actorDID string) (int64, error) { - // Consent/kill-switch withdrawal: park the actor's PENDING work as - // cancelled (never poisoned — this is not a failure). Terminal deliveries - // are left untouched. Joined through outbound_activities.actor_did. + return cancelForActor(ctx, r.db, actorDID) +} + +func (r *postgresOutboundDeliveries) CancelForActorTx(ctx context.Context, tx *sql.Tx, actorDID string) (int64, error) { + if tx == nil { + return 0, errors.NewValidationError("tx", "must not be nil") + } + return cancelForActor(ctx, tx, actorDID) +} + +// 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. +func cancelForActor(ctx context.Context, ex execer, actorDID string) (int64, error) { query := ` UPDATE outbound_deliveries d SET state = 'cancelled', claimed_until = NULL, updated_at = now() @@ -236,7 +288,15 @@ func (r *postgresOutboundDeliveries) CancelForActor(ctx context.Context, actorDI AND a.actor_did = $1 AND d.state = 'pending'` - return r.cancel(ctx, "cancel outbound_deliveries for actor", query, actorDID) + result, err := ex.ExecContext(ctx, query, actorDID) + if err != nil { + return 0, fmt.Errorf("cancel outbound_deliveries for actor %q: %w", actorDID, err) + } + affected, err := result.RowsAffected() + if err != nil { + return 0, fmt.Errorf("cancel outbound_deliveries for actor %q: rows affected: %w", actorDID, err) + } + return affected, nil } func (r *postgresOutboundDeliveries) CancelForCommunity(ctx context.Context, orderingKey string) (int64, error) { @@ -300,6 +360,57 @@ func (r *postgresOutboundDeliveries) cancel(ctx context.Context, op, query strin return affected, nil } +func (r *postgresOutboundDeliveries) DistinctInboxesForActor(ctx context.Context, actorDID string) ([]DeliveryTarget, error) { + if actorDID == "" { + return nil, errors.NewValidationError("actor_did", "must not be empty") + } + // THE DELIVERY HISTORY IS THE ADDRESS BOOK. There is no other record of + // which instances hold a user's content: an erasure has to go where the + // content actually went, and that is these rows. + // + // DISTINCT ON the inbox, because the fan-out is per INSTANCE: co-hosted + // communities share one inbox, and asking it twice sends the same instance + // the same erasure twice. Each surviving row keeps a REAL ordering key from + // the history, so the withdrawal serializes on a line the actor's other work + // already uses rather than jumping an independent queue. + // + // EVERY state counts, terminal ones included. A delivered post is exactly + // what has to be withdrawn; a poisoned or cancelled one may still have + // reached the peer (the wire and the ledger disagree by definition in those + // states), and reaching an instance that does not hold the content is + // harmless while missing one that does is not. + // + // KNOWN GAP, and it is the honest boundary of what this can reach: an + // instance that discovered the actor through search or WebFinger and never + // received a delivery from us has no row here, so the erasure never reaches + // it. Nothing in the bridge's state knows about that instance. + query := ` + SELECT DISTINCT ON (d.target_inbox) d.target_inbox, d.ordering_key + FROM outbound_deliveries d + JOIN outbound_activities a ON a.activity_id = d.activity_id + WHERE a.actor_did = $1 + ORDER BY d.target_inbox, d.seq` + + rows, err := r.db.QueryContext(ctx, query, actorDID) + if err != nil { + return nil, fmt.Errorf("list delivery inboxes for %q: %w", actorDID, err) + } + defer func() { _ = rows.Close() }() + + var targets []DeliveryTarget + for rows.Next() { + var target DeliveryTarget + if err := rows.Scan(&target.Inbox, &target.OrderingKey); err != nil { + return nil, fmt.Errorf("scan delivery inbox for %q: %w", actorDID, err) + } + targets = append(targets, target) + } + if err := rows.Err(); err != nil { + return nil, fmt.Errorf("list delivery inboxes for %q: %w", actorDID, err) + } + return targets, nil +} + func (r *postgresOutboundDeliveries) Get(ctx context.Context, activityID, targetInbox string) (*OutboundDelivery, error) { query := `SELECT` + deliveryColumns + ` FROM outbound_deliveries WHERE activity_id = $1 AND target_inbox = $2` diff --git a/internal/store/outbound_votes.go b/internal/store/outbound_votes.go index f7c6185..e676887 100644 --- a/internal/store/outbound_votes.go +++ b/internal/store/outbound_votes.go @@ -125,6 +125,42 @@ func (r *postgresOutboundVotes) GetByActivityID(ctx context.Context, activityID return vote, nil } +func (r *postgresOutboundVotes) ListDeliveredForActor(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. + // + // Served by the partial index on the same predicate (migration 029). + query := `SELECT` + outboundVoteColumns + ` + FROM outbound_votes + WHERE actor_did = $1 AND delivered_state = 'delivered' + ORDER BY 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) + } + defer func() { _ = rows.Close() }() + + var votes []OutboundVote + for rows.Next() { + vote, err := scanOutboundVote(rows) + if err != nil { + return nil, fmt.Errorf("scan delivered 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 votes, nil +} + func (r *postgresOutboundVotes) SetDeliveredState(ctx context.Context, voteATURI string, state DeliveredState) error { // Validated in Go rather than left to the CHECK constraint: an unknown // state is a caller bug, and the caller needs it back as a validation