diff --git a/internal/consume/account.go b/internal/consume/account.go index 0b118a5..fa9b4ff 100644 --- a/internal/consume/account.go +++ b/internal/consume/account.go @@ -5,47 +5,186 @@ import ( "database/sql" stderrors "errors" "fmt" + "hash/fnv" + "log/slog" ) // The #account ordering guard (second-opinion C4). #account frames are UNGATED // by the rev gate — they carry no rev — so their monotonic per-DID seq is the // only thing standing between a reconnect rewind and a silently regressed // account state. ap_actors.last_account_seq (migration 019) holds the -// high-water mark per actor, and these two helpers are the read and the -// monotonic advance that bracket the handler, exactly like jetstream_record_revs -// brackets a commit. +// high-water mark per actor, and this file is the CLAIM that brackets the +// handler, exactly like applyGatedTx brackets a commit. // // The seq lives on ap_actors because the gate only ever runs for a DID that // already has an actor: the actor-existence check precedes it, so an actorless // account event is skipped before any seq is consulted or recorded. +// +// WHY A CLAIM AND NOT A READ. main runs connector.Start and the DLQ redriver as +// two goroutines over ONE Dispatcher, so a redriven frame and a live frame for +// the same DID are handled concurrently. Read-decide-write-advance as four +// autocommit statements loses both races that matters here: two racers read the +// same watermark, both pass the gate, and the STALE one's write can land last +// (a reactivated user silently left paused until the next live frame, which may +// be days away) — or, on status="deleted", both reach the terminal, irreversible +// seam and two Delete{Person} activities go out for one deletion. + +// accountClaimNamespace scopes this package's advisory locks. The two-argument +// advisory lock functions take (namespace, key) as a pair of int32s and occupy a +// lock space DISTINCT from the one-argument int64 form — which is what keeps +// this from colliding with internal/repo's commit lock or testutil's +// process-wide harness lock, both of which use the one-argument form. +const accountClaimNamespace = 0x6163 // "ac" + +// accountClaimKey is the per-DID half of the advisory lock key. +// +// Hashed in GO rather than with postgres's hashtext(): that function is +// undocumented and its algorithm has changed across major versions, and a lock +// key is not something to inherit from an implementation detail. FNV-1a is +// fixed here forever, which is all a lock key has to be. +// +// A COLLISION IS SAFE, not merely unlikely: two DIDs sharing a key serialize +// against each other, which costs a little contention on the rarest event kind +// the consumer handles and changes no outcome. Correctness needs only that one +// DID always maps to one key. +func accountClaimKey(did string) int32 { + hash := fnv.New32a() + _, _ = hash.Write([]byte(did)) // hash.Write never returns an error + return int32(hash.Sum32()) //nolint:gosec // wraparound is fine: this is a lock key, not a value +} + +// errAccountUnapplied is the sentinel an apply function returns to ABANDON its +// claim without failing the event: the transition did not happen, so the seq +// must not advance, but there is nothing to retry either (an unwired terminal +// tier, an actor that vanished mid-flight). The claim rolls back and the event +// is reported handled, which leaves it replayable — the rollback is the whole +// point, and it must not be mistaken for a storage failure. +var errAccountUnapplied = stderrors.New("account event left unapplied") + +// applyAccountClaimed runs apply under an exclusive per-DID claim on the +// #account seq, and advances the seq only if apply succeeds. It is +// applyGatedTx's shape — claim first, apply under it, commit last — with ONE +// deliberate divergence, and the divergence is the whole reason this is a +// separate function rather than a call into the rev gate: +// +// THE CLAIM IS AN ADVISORY LOCK, NOT THE ap_actors ROW LOCK. The obvious +// implementation is `UPDATE ap_actors SET last_account_seq = $2 WHERE did = $1 +// AND last_account_seq < $2` as the first statement, holding that row's lock +// across apply. It deadlocks. The deleted path's apply calls the terminal tier, +// which takes NO transaction and opens its own — and the destructive tier it +// reaches (outbound.Purger) tombstones the very ap_actors row this transaction +// would be holding. That is exactly the hazard rev_gate.go's DEADLOCK NOTE +// names, and the reason handleFederation defers its own ap_actors write until +// after the destructive seam returns: this transaction cannot release what it +// holds without committing, and it is synchronously waiting on the call that +// needs it. +// +// An advisory lock excludes the other #account handlers — the only writers of +// last_account_seq — while contending with nothing the terminal tier touches. +// The row is written LAST, so its lock is held for the commit and not a +// microsecond longer. +// +// THE CLAIM IS HELD ACROSS THE TERMINAL TIER'S NETWORK PROBE, on purpose. That +// tier confirms the deletion against PLC and the PDS before sending anything, +// and the alternative shapes are worse: releasing the claim around the probe +// re-opens precisely the window in which two racers both confirm and both send +// Delete{Person}, and claiming BEFORE the probe (advancing the seq, then +// probing) consumes the event so a failed confirm can never be retried — the +// case account_confirm_test.go pins as the reason the seam exists. The cost is a +// pooled connection held for the length of a bounded HTTP round trip on the +// rarest event kind the consumer handles. +// +// A claim LOSS returns nil: the event is a duplicate or a stale copy, fully +// accounted for, and the cursor must advance past it. +func (d *Dispatcher) applyAccountClaimed(ctx context.Context, did string, seq int64, apply func(tx *sql.Tx) error) error { + tx, err := d.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin account claim for %s: %w", did, err) + } + defer func() { + if rollbackErr := tx.Rollback(); rollbackErr != nil && !stderrors.Is(rollbackErr, sql.ErrTxDone) { + d.logger.Error("failed to roll back account claim", + slog.String("did", did), slog.String("error", rollbackErr.Error())) + } + }() + + // FIRST STATEMENT. Everything below — the watermark read, the decision, the + // transition, the advance — happens with every other #account handler for + // this DID shut out, which is what makes the read-then-write below a claim + // rather than a guess. + if _, err := tx.ExecContext(ctx, + `SELECT pg_advisory_xact_lock($1, $2)`, + accountClaimNamespace, accountClaimKey(did)); err != nil { + return fmt.Errorf("claim account seq for %s: %w", did, err) + } -// lastAccountSeq returns the highest #account seq applied for the DID's actor. -// A DID with no actor row reports 0 — but the caller reaches this only after -// confirming the actor exists, so in practice the row is always present. -func (d *Dispatcher) lastAccountSeq(ctx context.Context, did string) (int64, error) { - var seq int64 - err := d.db.QueryRowContext(ctx, - `SELECT last_account_seq FROM ap_actors WHERE did = $1`, did).Scan(&seq) + var lastSeq int64 + err = tx.QueryRowContext(ctx, + `SELECT last_account_seq FROM ap_actors WHERE did = $1`, did).Scan(&lastSeq) if stderrors.Is(err, sql.ErrNoRows) { - return 0, nil + // The actor vanished between the existence check and the claim. Nothing + // to pause and nothing to withdraw; no seq is recorded, so a redelivery + // still applies if the actor is minted again. + d.logger.Debug("account claim found no actor", slog.String("did", did)) + return nil } if err != nil { - return 0, fmt.Errorf("read account seq for %s: %w", did, err) + return fmt.Errorf("read account seq for %s: %w", did, err) + } + if seq <= lastSeq { + d.logger.Debug("skipping stale or duplicate #account", + slog.String("did", did), slog.Int64("seq", seq), slog.Int64("last_seq", lastSeq)) + return nil + } + + if err := apply(tx); err != nil { + if stderrors.Is(err, errAccountUnapplied) { + // Deliberately un-advanced: the deferred rollback releases the claim + // with the watermark where it was, so the event can still be applied + // by a build (or a moment) that can act on it. + return nil + } + return err // the deferred rollback releases the claim un-advanced } - return seq, nil -} -// advanceAccountSeq records seq as the last applied #account seq for the DID's -// actor. The WHERE clause keeps it MONOTONIC: a late or out-of-order write can -// only ever raise the mark, never lower it, so an advance that races a newer -// one cannot reopen the door to a stale replay. -func (d *Dispatcher) advanceAccountSeq(ctx context.Context, did string, seq int64) error { - _, err := d.db.ExecContext(ctx, ` + // The advance is the LAST write, and it keeps its own monotonic WHERE. The + // claim already guarantees no concurrent advance; the predicate costs + // nothing and keeps the statement correct on its own terms. + if _, err := tx.ExecContext(ctx, ` UPDATE ap_actors SET last_account_seq = $2, updated_at = now() - WHERE did = $1 AND last_account_seq < $2`, - did, seq) - if err != nil { + WHERE did = $1 AND last_account_seq < $2`, did, seq); err != nil { return fmt.Errorf("advance account seq for %s: %w", did, err) } + if err := tx.Commit(); err != nil { + return fmt.Errorf("commit account claim for %s: %w", did, err) + } + return nil +} + +// setActorPausedTx writes the transient #account state ON THE CLAIM's +// transaction, so the pause and the seq advance land together or not at all. +// +// It is raw SQL here rather than a store.APActors method for the same reason +// the seq itself is: this is the consume package's own gate surface, the +// statement is the mirror of store's SetPaused (delivery_paused alone — pausing +// delivery must never disable the actor), and the claim needs it on a +// transaction the store's interface does not offer one for. +// +// A missing actor reports errAccountUnapplied: the row disappeared under us, so +// nothing was applied and the seq must not move. +func setActorPausedTx(ctx context.Context, tx *sql.Tx, did string, paused bool) error { + result, err := tx.ExecContext(ctx, + `UPDATE ap_actors SET delivery_paused = $2, updated_at = now() WHERE did = $1`, + did, paused) + if err != nil { + return fmt.Errorf("set delivery paused for %s: %w", did, err) + } + affected, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("set delivery paused for %s: rows affected: %w", did, err) + } + if affected == 0 { + return fmt.Errorf("%w: no actor for %s", errAccountUnapplied, did) + } return nil } diff --git a/internal/consume/account_claim_test.go b/internal/consume/account_claim_test.go new file mode 100644 index 0000000..7672b43 --- /dev/null +++ b/internal/consume/account_claim_test.go @@ -0,0 +1,372 @@ +package consume + +import ( + "context" + "database/sql" + "fmt" + "io" + "log/slog" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/optout" + "tidepool/internal/store" +) + +// THE #account TIER IS A CHECK-THEN-ACT, AND TWO GOROUTINES RUN IT. +// +// main starts connector.Start and the DLQ redriver over the SAME Dispatcher, so +// a redriven frame and a live frame are handled CONCURRENTLY. Commit events are +// safe: applyGatedTx claims the gate row as the first statement of a +// transaction and holds it across apply, so a racer blocks on the claim and +// then observes it. #account had none of that — read last_account_seq, decide, +// write, advance, four separate autocommit statements — which is two distinct +// bugs at once: +// +// 1. a stale pause can be applied AFTER a newer reactivation, because both +// racers read the same watermark before either advanced it, leaving the +// user un-delivered until the next live event; and +// 2. status="deleted" reaches the TERMINAL, IRREVERSIBLE seam TWICE, which is +// two Delete{Person} activities for one deletion. +// +// The claim these tests demand is applyGatedTx's shape: exclusive first, apply +// under it, commit last. + +// blockingTerminator is the terminal seam held OPEN. It is what makes the race +// deterministic rather than probabilistic: while one racer is parked inside the +// irreversible call, the other has all the time it needs to reach its own +// decision — which is exactly the window a claim has to close. +type blockingTerminator struct { + mu sync.Mutex + calls []string + entered chan string + release chan struct{} +} + +func newBlockingTerminator() *blockingTerminator { + return &blockingTerminator{ + entered: make(chan string, 4), + release: make(chan struct{}), + } +} + +func (t *blockingTerminator) TerminateAccount(_ context.Context, did string) error { + t.mu.Lock() + t.calls = append(t.calls, did) + t.mu.Unlock() + t.entered <- did + <-t.release + return nil +} + +func (t *blockingTerminator) Calls() []string { + t.mu.Lock() + defer t.mu.Unlock() + return append([]string(nil), t.calls...) +} + +// handleConcurrently dispatches every frame at once and returns their errors in +// order. Each frame is parsed on THIS goroutine (a redriver parses its own copy +// of the frame anyway) so no testing.T call happens off the test goroutine. +func handleConcurrently(t *testing.T, dispatcher *Dispatcher, frames ...[]byte) []error { + t.Helper() + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + defer cancel() + + events := make([]*JetstreamEvent, len(frames)) + for i, frame := range frames { + events[i] = parseFrame(t, frame) + } + errs := make([]error, len(frames)) + start := make(chan struct{}) + var wg sync.WaitGroup + for i := range events { + wg.Add(1) + go func(i int) { + defer wg.Done() + <-start + errs[i] = dispatcher.HandleEvent(ctx, events[i]) + }(i) + } + close(start) // released together, so the check-then-act window is real + wg.Wait() + return errs +} + +// --------------------------------------------------------------------------- +// Finding 1 — the claim +// --------------------------------------------------------------------------- + +// TestAccountClaim_ConcurrentDeletedReplayTerminatesExactlyOnce is the probe +// from the review, as a test. A redriven deletion and its live twin arrive +// together; without a claim BOTH pass the seq gate and BOTH call the terminal +// tier, and the terminal tier is the one that cannot be taken back. +func TestAccountClaim_ConcurrentDeletedReplayTerminatesExactlyOnce(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + terminator := newBlockingTerminator() + fixture := newDispatchFixture(t, database, + func(opts *Options) { opts.Terminator = terminator }) + + var errs []error + done := make(chan struct{}) + go func() { + defer close(done) + errs = handleConcurrently(t, fixture.dispatcher, + accountFrameSeq(dispatchNativeDID, false, "deleted", 41), + accountFrameSeq(dispatchNativeDID, false, "deleted", 41)) + }() + + select { + case <-terminator.entered: + case <-time.After(30 * time.Second): + t.Fatal("no racer reached the terminal seam") + } + // The other racer now has an unhurried window to read the watermark, decide + // and act. Under a claim it spends that window BLOCKED on the claim; without + // one it spends it sending a second Delete{Person}. + time.Sleep(250 * time.Millisecond) + close(terminator.release) + <-done + + for _, err := range errs { + require.NoError(t, err) + } + assert.Len(t, terminator.Calls(), 1, + "ONE deletion, ONE withdrawal. The terminal tier is irreversible at every peer "+ + "at once, so a redriven frame racing its live twin must be serialized by a "+ + "claim — not by hoping the two reads do not overlap") + assert.Equal(t, int64(41), appliedSeq(t, database, dispatchNativeDID)) +} + +// TestAccountClaim_StalePauseCannotLandAfterAReactivation is the other half of +// the same missing claim: the seq gate is read-then-write, so two racers can +// both pass it and the LOSER's write can land last. +func TestAccountClaim_StalePauseCannotLandAfterAReactivation(t *testing.T) { + database := dispatchTestDB(t) + fixture := newDispatchFixture(t, database) + + // Rounds, because an unguarded interleaving is a race rather than a + // certainty: each round is a fresh actor, and ONE bad round is the bug. + const rounds = 25 + for round := range rounds { + did := fmt.Sprintf("did:plc:racer%019d", round) + seedAPActor(t, database, did, fmt.Sprintf("racer%d", round)) + + errs := handleConcurrently(t, fixture.dispatcher, + accountFrameSeq(did, false, "deactivated", 5), // the redriven stale pause + accountFrameSeq(did, true, "active", 6), // the live reactivation + ) + for _, err := range errs { + require.NoError(t, err, "round %d", round) + } + + require.False(t, deliveryPaused(t, database, did), + "round %d: the seq-5 pause is STALE — seq 6 says the user is back. Applied in "+ + "either order the outcome must be the newer state, because the only thing "+ + "that repairs a wrongly paused actor is the NEXT live #account frame, and "+ + "there may not be one for days", round) + require.Equal(t, int64(6), appliedSeq(t, database, did), "round %d", round) + } +} + +// --------------------------------------------------------------------------- +// Finding 2 — a nil terminator must not eat the deletion +// --------------------------------------------------------------------------- + +// TestAccountDeleted_NilTerminalTierLeavesTheEventReplayable pins the one +// asymmetry that makes an unwired seam unrecoverable rather than merely +// unimplemented. Advancing the seq CLOSES the event: every redelivery of it is +// then rejected as stale, so a deployment that wires the terminal tier tomorrow +// can never act on the deletion that arrived today. +// +// The sibling RemoteContentDeleter path gets this right by making the durable +// record (the preference) BEFORE it reaches the seam. This path has nothing to +// record, so the un-advanced seq IS the record. +func TestAccountDeleted_NilTerminalTierLeavesTheEventReplayable(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + unwired := newDispatchFixture(t, database, func(opts *Options) { opts.Terminator = nil }) + + require.NoError(t, unwired.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 51)), + "a missing terminal tier must not fail the event") + assert.Zero(t, appliedSeq(t, database, dispatchNativeDID), + "and must not CONSUME it either: nothing was recorded and nothing was sent, so "+ + "advancing the seq would make the user's deletion permanently unreachable — "+ + "every later redelivery rejected as stale, by a bridge that never acted on it") + + // The deployment that wires the tier gets the event. + terminator := &recordingTerminator{} + wired := newDispatchFixture(t, database, + func(opts *Options) { opts.Terminator = terminator }) + require.NoError(t, wired.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 51))) + assert.Equal(t, []string{dispatchNativeDID}, terminator.DIDs(), + "the same frame, redelivered to a build with the tier wired, still reaches it") + assert.Equal(t, int64(51), appliedSeq(t, database, dispatchNativeDID), + "and NOW the event is closed") +} + +// --------------------------------------------------------------------------- +// Finding 3 — the terminal marker +// --------------------------------------------------------------------------- + +// actorTombstoned reports whether the destructive tier withdrew the identity. +func actorTombstoned(t *testing.T, database *sql.DB, did string) bool { + t.Helper() + var tombstoned bool + require.NoError(t, database.QueryRowContext(context.Background(), + `SELECT tombstoned_at IS NOT NULL FROM ap_actors WHERE did = $1`, did).Scan(&tombstoned)) + return tombstoned +} + +// TestAccountDeleted_TerminalStateReachesTheActorAndSurvivesAReactivation is +// the deployment where the DESTRUCTIVE seam has not landed: the terminal tier +// confirms the deletion and RECORDS it, and that record is the only terminal +// state in the database — no tombstone, because nothing was sent. +// +// The consume side wrote nothing at all: consent stayed OK on the actor row, +// which then disagreed with the preference that says this user is gone. Every +// reader that consults ap_actors WITHOUT the preference — webfinger, the actor +// document, the admission gate — went on treating a confirmed-deleted identity +// as live, and a later active=true frame with a greater seq is welcomed by a row +// that never learned anything happened. +func TestAccountDeleted_TerminalStateReachesTheActorAndSurvivesAReactivation(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + world := newTerminalWorld(t, database, false) // no destructive seam wired + world.confirmer.set(true, nil) + + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 61))) + + pref := storedPref(t, database, dispatchNativeDID) + require.NotNil(t, pref, "precondition: the terminal tier recorded the deletion") + require.False(t, pref.Enabled) + + enabled, _ := actorEnabled(t, database, dispatchNativeDID) + assert.False(t, enabled, + "THE ACTOR ROW MUST AGREE WITH THE PREFERENCE. The opt-out door mirrors exactly "+ + "this onto ap_actors; the deletion door — the stronger of the two — left the "+ + "identity enabled, so it kept resolving through webfinger for a user the "+ + "bridge had just confirmed is gone") + + // The reactivation that must find a closed door. + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, true, "active", 62))) + + enabled, _ = actorEnabled(t, database, dispatchNativeDID) + assert.False(t, enabled, + "an active=true frame at a greater seq must NOT re-open an identity whose deletion "+ + "was already acted on: the transient tier pauses and unpauses delivery, and it "+ + "has no business undoing a terminal decision") + assert.False(t, storedPref(t, database, dispatchNativeDID).Enabled, + "and the terminal record still stands") + assert.Equal(t, int64(62), appliedSeq(t, database, dispatchNativeDID)) +} + +// tombstoningDeleter is the destructive seam modelled on the REAL one: an +// outbound.Purger's last act is ap_actors.TombstoneTx, and that tombstone — +// not anything the consumer writes — is where a purged deployment's terminality +// actually lives. The double writes it so this package can test what the +// consume side must NOT undo, without depending on outbound. +type tombstoningDeleter struct { + actors store.APActors +} + +func (d *tombstoningDeleter) DeleteRemoteContent(ctx context.Context, did string) error { + return d.actors.Tombstone(ctx, did) +} + +// TestAccountDeleted_PurgedIdentityIsNotResurrectedByAReactivation is the same +// question with the destructive tier WIRED, where terminality comes from the +// purge (tombstoned_at, which no code path clears). It pins that the consume +// side's own marker COMPOSES with the purge's rather than fighting it. +func TestAccountDeleted_PurgedIdentityIsNotResurrectedByAReactivation(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + confirmer := &stubConfirmer{} + confirmer.set(true, nil) + terminator, err := optout.NewTerminator(optout.Options{ + Confirmer: confirmer, + Prefs: store.NewFederationPrefs(database), + Deleter: &tombstoningDeleter{actors: store.NewAPActors(database)}, + Logger: slog.New(slog.NewTextHandler(io.Discard, nil)), + }) + require.NoError(t, err) + world := newDispatchFixture(t, database, func(opts *Options) { opts.Terminator = terminator }) + + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, false, "deleted", 71))) + require.True(t, actorTombstoned(t, database, dispatchNativeDID), + "precondition: the purge withdrew the identity") + + require.NoError(t, world.handle(t, accountFrameSeq(dispatchNativeDID, true, "active", 72))) + + assert.True(t, actorTombstoned(t, database, dispatchNativeDID), + "a withdrawal is terminal: peers were told this identity is gone and no peer "+ + "un-deletes, so nothing an #account frame says may clear it") + enabled, _ := actorEnabled(t, database, dispatchNativeDID) + assert.False(t, enabled, + "and it stays disabled — signing new content as an actor peers were told is gone "+ + "is the one outcome the tombstone exists to prevent") +} + +// --------------------------------------------------------------------------- +// Finding 4 — a frame with no seq +// --------------------------------------------------------------------------- + +// TestAccount_FrameWithoutAPositiveSeqIsPermanentlyRejected closes the silent +// hole under the seq gate. last_account_seq defaults to 0 and the gate skips +// anything at or below it, so a frame carrying seq 0 — or none at all, which +// decodes to the same thing — is indistinguishable from a stale replay ON THE +// VERY FIRST EVENT. A Jetstream build or proxy that stopped emitting seq would +// take the whole account tier dark with clean metrics and a debug line reading +// "stale or duplicate". +// +// time_us gets exactly this treatment for every kind, for exactly this reason. +func TestAccount_FrameWithoutAPositiveSeqIsPermanentlyRejected(t *testing.T) { + for _, tc := range []struct { + name string + frame []byte + }{ + {"zero", []byte(fmt.Sprintf( + `{"did":%q,"time_us":7300,"kind":"account",`+ + `"account":{"did":%q,"seq":0,"time":"2026-08-13T10:00:00.000Z",`+ + `"active":false,"status":"deactivated"}}`, + dispatchNativeDID, dispatchNativeDID))}, + {"absent", []byte(fmt.Sprintf( + `{"did":%q,"time_us":7301,"kind":"account",`+ + `"account":{"did":%q,"time":"2026-08-13T10:00:00.000Z",`+ + `"active":false,"status":"deactivated"}}`, + dispatchNativeDID, dispatchNativeDID))}, + {"negative", []byte(fmt.Sprintf( + `{"did":%q,"time_us":7302,"kind":"account",`+ + `"account":{"did":%q,"seq":-3,"time":"2026-08-13T10:00:00.000Z",`+ + `"active":false,"status":"deactivated"}}`, + dispatchNativeDID, dispatchNativeDID))}, + } { + t.Run(tc.name, func(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + terminator := &recordingTerminator{} + fixture := newDispatchFixture(t, database, + func(opts *Options) { opts.Terminator = terminator }) + + err := fixture.handle(t, tc.frame) + require.Error(t, err, + "a frame with no usable seq cannot be ORDERED, and an #account tier that "+ + "cannot order is not a tier — it is a coin flip between the stale copy "+ + "and the live one") + assert.ErrorIs(t, err, ErrPermanentEvent, + "and no retry adds a seq to a frame that has none, so it is dead-lettered "+ + "where an operator can SEE the tier went dark instead of reading a debug "+ + "line that says 'stale or duplicate'") + + assert.False(t, deliveryPaused(t, database, dispatchNativeDID), + "nothing is applied from a frame that cannot be ordered") + assert.Zero(t, appliedSeq(t, database, dispatchNativeDID)) + assert.Empty(t, terminator.DIDs()) + }) + } +} diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index 0f932be..ad6f1f5 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -594,6 +594,20 @@ func (d *Dispatcher) handleAccount(ctx context.Context, event *JetstreamEvent) e } did := event.DID + // seq is the ORDERING position, and it is validated here for the same + // reason time_us is validated for every kind above. last_account_seq + // defaults to 0 and the gate below skips anything at or below it, so a + // frame carrying seq 0 — or none at all, which decodes to the same zero — + // is indistinguishable from a stale replay ON THE VERY FIRST EVENT: a + // Jetstream build or proxy that stopped emitting seq would take the whole + // account tier dark, every frame logged as "stale or duplicate", metrics + // clean. No retry adds a seq to a frame that has none, so it is dead- + // lettered where an operator can see it instead of being skipped in silence. + if account.Seq <= 0 { + return fmt.Errorf("%w: account event for %s carries no positive seq, got %d", + ErrPermanentEvent, strconv.Quote(did), account.Seq) + } + // The actor check comes first for EVERY status. Nothing was ever federated // under a DID with no AP identity, so there is no delivery to pause and // nothing for the terminal tier to withdraw — and minting an actor in @@ -614,16 +628,25 @@ func (d *Dispatcher) handleAccount(ctx context.Context, event *JetstreamEvent) e // reconnect rewind (a pause at seq N redelivered after a reactivation at // N+1) must not flip a recovered user back. A seq at or below the last // applied one is a duplicate or a stale copy and is a no-op. - lastSeq, err := d.lastAccountSeq(ctx, did) - if err != nil { - return err - } - if account.Seq <= lastSeq { - d.logger.Debug("skipping stale or duplicate #account", - slog.String("did", did), slog.Int64("seq", account.Seq), slog.Int64("last_seq", lastSeq)) - return nil - } + // + // The gate is CLAIMED rather than read (see applyAccountClaimed): the + // redriver and the live connector share this dispatcher, so a read-then-act + // gate lets both racers through — and on the deleted path both then reach a + // seam that no peer can undo. + return d.applyAccountClaimed(ctx, did, account.Seq, func(tx *sql.Tx) error { + return d.applyAccountStatus(ctx, tx, did, account) + }) +} +// applyAccountStatus performs the transition for one #account frame. It runs +// under the claim, so it may assume this is the newest frame seen for the DID +// and that no other #account handler is running for it. +// +// ORDERING RULE, and it is the same one rev_gate.go's DEADLOCK NOTE states for +// commit handlers: the terminal tier opens its OWN transaction and writes the +// ap_actors row, so nothing here may write that row BEFORE the seam returns. +// Every write below happens after it. +func (d *Dispatcher) applyAccountStatus(ctx context.Context, tx *sql.Tx, did string, account *AccountEvent) error { if account.Status == accountStatusDeleted { // The ONE status that means gone. It goes to the tier that re-verifies // against PLC and the PDS before sending Delete{Person}, because @@ -632,28 +655,72 @@ func (d *Dispatcher) handleAccount(ctx context.Context, event *JetstreamEvent) e // Announced, and deliberately NOT degraded into a pause: a // deletion half-handled as a pause looks handled in the database // and is not. - d.logger.Warn("account reported deleted but no terminal tier is wired", - slog.String("did", did)) - return d.advanceAccountSeq(ctx, did, account.Seq) + // + // THE SEQ DOES NOT ADVANCE. Nothing was recorded and nothing was + // sent, so advancing it would CONSUME the user's deletion: every + // redelivery is then rejected as stale and a build that wires the + // tier tomorrow can never act on the event that arrived today. The + // sibling deleteRemote path can advance because it persists the + // preference before reaching its seam; here the un-advanced seq is + // the only record that the deletion is still owed. + d.logger.Warn("account reported deleted but no terminal tier is wired; "+ + "the event is left unapplied so a build with the tier can still act on it", + slog.String("did", did), slog.Int64("seq", account.Seq)) + return errAccountUnapplied } if err := d.terminator.TerminateAccount(ctx, did); err != nil { return fmt.Errorf("terminate account %s: %w", did, err) } - return d.advanceAccountSeq(ctx, did, account.Seq) + return d.mirrorTerminalPreference(ctx, tx, did) } // Everything else is transient — deactivated, suspended, takendown, // throttled are all states a user comes back from. Delivery stops; the // identity, and every federated reference to it, survives. - if err := d.apActors.SetPaused(ctx, did, !account.Active); err != nil { - if errors.IsNotFound(err) { - return nil // the actor vanished between the check and the write - } - return err + // + // Note what this does NOT do: it never re-ENABLES an actor. Unpausing + // restores delivery for an identity that still exists, while `enabled` + // carries the terminal decisions — an opt-out, a confirmed deletion, a + // withdrawal — and an active=true frame is not evidence that any of those + // were reversed. A tombstoned actor's store refuses re-enabling outright; + // this path simply never asks. + return setActorPausedTx(ctx, tx, did, !account.Active) +} + +// mirrorTerminalPreference brings the ap_actors mirror into agreement with +// federation_prefs after the terminal tier has run — the SAME invariant the +// opt-out door maintains (handleFederation writes the preference, then mirrors +// it onto the actor), applied at the deletion door, which never did. +// +// IT READS RATHER THAN DECIDES, because TerminateAccount returns nil for two +// opposite outcomes: a CONFIRMED deletion (which records enabled=false and, if +// the destructive seam is wired, tombstones the actor) and a confirmed LIVE +// account (which withdraws the tier's own stale request, leaving absence, which +// means default-on). The preference is the tier's own statement of which +// happened, so reading it is what keeps this from re-deciding a question that +// was answered against PLC. +// +// It COMPOSES with the purge rather than duplicating it: a wired destructive +// tier has already stamped tombstoned_at — terminal, and cleared by nothing — +// and this write agrees with it. An UNWIRED one leaves the preference as the +// only record, and without this the actor row went on saying the identity is +// live: still resolving through webfinger, still admissible, and a later +// active=true frame arriving at a row that never learned anything happened. +func (d *Dispatcher) mirrorTerminalPreference(ctx context.Context, tx *sql.Tx, did string) error { + pref, err := d.prefs.Get(ctx, did) + if errors.IsNotFound(err) { + // Confirmed live: the tier withdrew its request and absence is + // default-on. Disabling here would strand a user the confirm just + // proved is still there. + return nil + } + if err != nil { + return fmt.Errorf("read federation preference for %s: %w", did, err) + } + if pref.Enabled { + return nil } - // Record the applied seq only AFTER the state change succeeds, so a failed - // write replays instead of being locked out by an advanced seq. - return d.advanceAccountSeq(ctx, did, account.Seq) + return d.mirrorActorEnabled(ctx, tx, did, false) } // accountStatusDeleted is the ONLY #account status that means deletion.