diff --git a/cmd/tidepool/main.go b/cmd/tidepool/main.go index ecc98fe..a478931 100644 --- a/cmd/tidepool/main.go +++ b/cmd/tidepool/main.go @@ -5,6 +5,7 @@ package main import ( "context" + "database/sql" "errors" "fmt" "log/slog" @@ -20,6 +21,7 @@ import ( "tidepool/internal/ap" "tidepool/internal/config" + "tidepool/internal/consume" "tidepool/internal/db" "tidepool/internal/identity" "tidepool/internal/ingest" @@ -462,6 +464,18 @@ func run(logger *slog.Logger) error { return fmt.Errorf("user origin: %w", err) } + // The Jetstream consumer (task 14): the atproto half of the world this + // bridge does not host. Default OFF until task 18 wires the e2e path — + // it writes durable outbound state, so a deployment that has not been + // wired end to end must not start accumulating it. + var consumerDone <-chan struct{} + if cfg.ConsumerEnabled { + consumerDone, err = startConsumer(ctx, cfg, database, repoManager, personasService, logger) + if err != nil { + return err + } + } + // Host routing wraps everything: the chi router keeps answering for the // bridge hostname and its bridged-handle subdomains, the user origin // answers for its own Host, and an unrecognized Host is refused with 421 @@ -530,6 +544,17 @@ func run(logger *slog.Logger) error { case <-shutdownCtx.Done(): logger.Warn("backfill drain timed out; abandoning in-flight run (resumable on restart)") } + // Wait for the consumer's read loop to exit. Its shutdown path flushes + // the cursor on a fresh context, so cutting the process short here + // would lose the progress since the last periodic flush and replay it + // on the next boot. + if consumerDone != nil { + select { + case <-consumerDone: + case <-shutdownCtx.Done(): + logger.Warn("jetstream consumer did not stop in time; its cursor may replay on restart") + } + } // ListenAndServe has returned by now (Shutdown guarantees it); // drain its error so a bind failure racing the signal still exits // non-zero instead of being lost in the buffered channel. @@ -540,3 +565,93 @@ func run(logger *slog.Logger) error { return nil } } + +// startConsumer wires the Jetstream consumer (task 14) and starts its read +// loop. The returned channel closes when the connector's loop has exited, so +// shutdown can wait for the final cursor flush instead of racing it. +// +// The seams that are not wired yet are nil ON PURPOSE, and each is a no-op the +// consumer announces rather than a silent gap: +// +// - Enqueuer is the logging no-op until task 15's delivery queue lands. The +// consumer still runs behind it, so the cursor, the rev gate and the +// outbound state that delivery will be built FROM are all exercised. +// - Engine (task 16) nil means postv2 events are skipped at debug. +// - RemoteDeleter (task 17) nil means a deleteRemote opt-out is recorded and +// logged rather than acted on. +// - Terminator (task 17) nil means a deleted account is logged rather than +// withdrawn — never quietly downgraded to a delivery pause. +func startConsumer( + ctx context.Context, + cfg *config.Config, + database *sql.DB, + repoManager *repo.Manager, + minter consume.ActorMinter, + logger *slog.Logger, +) (<-chan struct{}, error) { + // The most SSRF-exposed egress in the bridge: the well-known host comes + // from a DID document a stranger controls, so it shares the AP client's + // guard rather than using a bare http.Client. + resolver, err := consume.NewHandleResolver(consume.ResolverOptions{ + PLCDirectoryURL: cfg.PLCDirectoryURL, + HTTPClient: ap.NewGuardedHTTPClient(cfg.AllowPrivateAddresses, 30*time.Second), + UserAgent: cfg.UserAgent, + // DNS is the first half of handle verification and covers every + // self-hosted handle that publishes no well-known. + LookupTXT: consume.DefaultLookupTXT, + Logger: logger, + }) + if err != nil { + return nil, fmt.Errorf("consumer: handle resolver: %w", err) + } + + dispatcher, err := consume.NewDispatcher(consume.Options{ + DB: database, + Actors: minter, + Resolver: resolver, + Enqueuer: consume.NewNoopEnqueuer(logger), + // Reads committed records so a subject's community resolves for + // mappings written before migration 016 filled community_did. + Records: repoManager, + UserOrigin: cfg.APUserOrigin, + Logger: logger, + }) + if err != nil { + return nil, fmt.Errorf("consumer: dispatcher: %w", err) + } + + // The collection filter is load-bearing: without wantedCollections this + // would subscribe to the entire network's firehose and discard it record + // by record. + subscribeURL, err := consume.SubscribeURL(cfg.JetstreamURL, consume.WantedCollections()) + if err != nil { + return nil, fmt.Errorf("consumer: %w", err) + } + + state := consume.NewPostgresStateStore(database, consume.CursorSchemaVersion) + connector := consume.NewConnector(consume.ConsumerNative, subscribeURL, dispatcher, + consume.WithCursorStore(state), + consume.WithDeadLetterWriter(state), + consume.WithConnectorLogger(logger)) + + // The redriver makes transient failures self-healing: an event captured + // during a postgres blip is replayed once the blip clears, without anyone + // being paged. + go consume.NewDeadLetterRedriver(state, + map[string]consume.EventHandler{consume.ConsumerNative: dispatcher}).Run(ctx) + + // Cursor age and dead-letter depth are what make a STALLED consumer + // visible: the process stays up and the health check stays green while + // events quietly stop arriving. + consume.PublishMetrics(context.Background(), connector, state) + + done := make(chan struct{}) + go func() { + defer close(done) + if err := connector.Start(ctx); err != nil { + logger.Error("jetstream consumer stopped", "error", err) + } + }() + logger.Info("jetstream consumer started", "url", subscribeURL) + return done, nil +} diff --git a/internal/config/config.go b/internal/config/config.go index 2d5ded3..d0f9c39 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -8,6 +8,7 @@ import ( "encoding/hex" "fmt" "log/slog" + "net/url" "os" "strconv" "strings" @@ -114,6 +115,17 @@ type Config struct { // refused: an authenticated write surface must not answer under a Host // an attacker chose. APHostFallthroughDev bool + + // ConsumerEnabled turns on the Jetstream consumer (task 14). Default OFF: + // the consumer writes durable outbound state and hands work to delivery, + // so a deployment that has not been wired end to end should not start + // silently accumulating it. + ConsumerEnabled bool + // JetstreamURL is the self-hosted Jetstream the consumer subscribes to + // (ws:// or wss://). REQUIRED when ConsumerEnabled — a consumer with + // nowhere to dial would come up "healthy" and consume nothing, which is + // the failure mode cursors and lag metrics exist to make impossible. + JetstreamURL string // AdminToken is the bearer token protecting the /admin API (community // subscribe/unsubscribe/backfill). ADMIN_TOKEN; required in production, // dev default is a fixed, publicly known value. @@ -406,6 +418,36 @@ func Load(logger *slog.Logger) (*Config, error) { return nil, err } + // The Jetstream consumer (task 14), default OFF until task 18 wires the + // e2e path: it writes durable outbound state and hands work to delivery, + // so a deployment that has not been wired end to end should not quietly + // start accumulating it. + cfg.ConsumerEnabled, err = boolVarDefault(logger, "CONSUMER_ENABLED", false) + if err != nil { + return nil, err + } + cfg.JetstreamURL = strings.TrimSpace(os.Getenv("JETSTREAM_URL")) + // Validated whenever it is SET, not only when the consumer is on: staging + // a URL ahead of the flag is how a deployment is prepared, and a typo + // caught then is a boot failure with a clear message instead of a + // reconnect loop on the day someone flips the switch. + if cfg.JetstreamURL != "" { + parsed, err := url.Parse(cfg.JetstreamURL) + if err != nil { + return nil, fmt.Errorf("config: JETSTREAM_URL is not a valid URL: %w", err) + } + if parsed.Scheme != "ws" && parsed.Scheme != "wss" || parsed.Host == "" { + return nil, fmt.Errorf("config: JETSTREAM_URL must be an absolute ws:// or wss:// URL, got %q", cfg.JetstreamURL) + } + } + if cfg.ConsumerEnabled && cfg.JetstreamURL == "" { + // A consumer with nowhere to dial comes up looking healthy and + // consumes NOTHING, and silence is indistinguishable from a quiet + // stream — the exact failure the cursor and lag metrics exist to + // expose. Refuse at boot instead. + return nil, fmt.Errorf("config: JETSTREAM_URL is required when CONSUMER_ENABLED is set") + } + // Retention knobs for the task-11 pruners: same semantics as // FIREHOSE_RETENTION (real defaults everywhere, must be positive). cfg.TombstoneRetention, err = durationVar(logger, "TOMBSTONE_RETENTION", 720*time.Hour) diff --git a/internal/config/config_test.go b/internal/config/config_test.go index 25b6cd1..29204d2 100644 --- a/internal/config/config_test.go +++ b/internal/config/config_test.go @@ -22,6 +22,7 @@ func clearConfigEnv(t *testing.T) { "MINT_BURST", "INGEST_WORKERS", "BRIDGE_SCHEME", "ALLOW_PRIVATE_FETCH", "ALLOW_DEV_REQUEST_CRAWL", "RELAY_HOSTS", "AP_USER_ORIGIN", "AP_HOST_FALLTHROUGH_DEV", + "CONSUMER_ENABLED", "JETSTREAM_URL", } { t.Setenv(name, "") } @@ -434,3 +435,92 @@ func TestLoad_APUserOriginShadowCheckIsCanonical(t *testing.T) { }) } } + +// --------------------------------------------------------------------------- +// Task 14 cycle K1: the Jetstream consumer's configuration. +// +// The consumer is default-OFF and stays that way until task 18 wires the e2e +// path, because it writes durable outbound state and hands work to delivery: +// a deployment that has not been wired end to end should not quietly start +// accumulating it. +// --------------------------------------------------------------------------- + +func TestLoad_ConsumerIsDisabledByDefault(t *testing.T) { + clearConfigEnv(t) + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + + assert.False(t, cfg.ConsumerEnabled, + "the consumer is off unless a deployment says otherwise") + assert.Empty(t, cfg.JetstreamURL, + "and an unset JETSTREAM_URL is fine while it is off") +} + +func TestLoad_EnabledConsumerRequiresAJetstreamURL(t *testing.T) { + clearConfigEnv(t) + t.Setenv("CONSUMER_ENABLED", "true") + + _, err := Load(discardLogger()) + require.Error(t, err, + "a consumer with nowhere to dial would come up looking healthy and consume "+ + "NOTHING — silence is indistinguishable from a quiet stream, which is "+ + "exactly the failure the cursor and lag metrics exist to expose") + assert.Contains(t, err.Error(), "JETSTREAM_URL", + "the message must name the variable an operator has to set") +} + +func TestLoad_EnabledConsumerAcceptsAWebSocketURL(t *testing.T) { + clearConfigEnv(t) + t.Setenv("CONSUMER_ENABLED", "1") + t.Setenv("JETSTREAM_URL", "ws://localhost:6008/subscribe") + + cfg, err := Load(discardLogger()) + require.NoError(t, err) + + assert.True(t, cfg.ConsumerEnabled) + assert.Equal(t, "ws://localhost:6008/subscribe", cfg.JetstreamURL) +} + +func TestLoad_JetstreamURLMustBeAWebSocketURL(t *testing.T) { + for _, raw := range []string{ + "https://jetstream.example/subscribe", // the scheme a copy-paste produces + "jetstream.example", // no scheme at all + "://nonsense", + } { + t.Run(raw, func(t *testing.T) { + clearConfigEnv(t) + t.Setenv("CONSUMER_ENABLED", "true") + t.Setenv("JETSTREAM_URL", raw) + + _, err := Load(discardLogger()) + require.Error(t, err, + "a bad URL must fail at BOOT: caught at dial time instead, it becomes a "+ + "reconnect loop that looks like an upstream outage") + assert.Contains(t, err.Error(), "JETSTREAM_URL") + }) + } +} + +func TestLoad_JetstreamURLIsCarriedEvenWhileDisabled(t *testing.T) { + clearConfigEnv(t) + t.Setenv("JETSTREAM_URL", "wss://jetstream.example/subscribe") + + cfg, err := Load(discardLogger()) + require.NoError(t, err, + "a URL configured ahead of the flag is not an error — that is how a deployment "+ + "is staged before being switched on") + assert.False(t, cfg.ConsumerEnabled) + assert.Equal(t, "wss://jetstream.example/subscribe", cfg.JetstreamURL) +} + +func TestLoad_ConsumerEnabledRejectsATypo(t *testing.T) { + clearConfigEnv(t) + t.Setenv("CONSUMER_ENABLED", "yess") + + _, err := Load(discardLogger()) + require.Error(t, err, + "boolVarDefault semantics: an unrecognised value is refused rather than read as "+ + "false, because a flag disabled by a typo is invisible") + assert.Contains(t, err.Error(), "CONSUMER_ENABLED") +} diff --git a/internal/consume/account_test.go b/internal/consume/account_test.go new file mode 100644 index 0000000..1149514 --- /dev/null +++ b/internal/consume/account_test.go @@ -0,0 +1,149 @@ +package consume + +import ( + "context" + "sync" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Task 14 cycle J: #account. +// +// Decision 19 in one sentence: active=false is NOT deletion. Deactivated, +// suspended, takendown and throttled are all states a user comes back from, +// and treating any of them as a deletion would send Delete{Person} to every +// peer — irreversibly destroying an identity over a temporary suspension. +// Only status="deleted" means gone, and even that goes to a seam that +// re-verifies before acting. + +// recordingTerminator stands in for the task 17 terminal tier. +type recordingTerminator struct { + mu sync.Mutex + dids []string + err error +} + +func (t *recordingTerminator) TerminateAccount(_ context.Context, did string) error { + t.mu.Lock() + defer t.mu.Unlock() + t.dids = append(t.dids, did) + return t.err +} + +func (t *recordingTerminator) DIDs() []string { + t.mu.Lock() + defer t.mu.Unlock() + return append([]string(nil), t.dids...) +} + +func TestAccount_TransientStatusesPauseDeliveryWithoutTouchingIdentity(t *testing.T) { + // Every inactive status atproto defines that is NOT a deletion. Each one + // is a user who may be back tomorrow. + for _, status := range []string{"deactivated", "suspended", "takendown", "throttled"} { + t.Run(status, func(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + terminator := &recordingTerminator{} + fixture := newDispatchFixture(t, database, + func(opts *Options) { opts.Terminator = terminator }) + + require.NoError(t, fixture.handle(t, + accountFrameFor(dispatchNativeDID, false, status))) + + assert.True(t, deliveryPaused(t, database, dispatchNativeDID), + "%s pauses DELIVERY", status) + enabled, _ := actorEnabled(t, database, dispatchNativeDID) + assert.True(t, enabled, + "but leaves the identity enabled: the actor document and every "+ + "federated reference to it must survive a temporary state") + assert.Empty(t, terminator.DIDs(), + "and NEVER reaches the terminal tier — sending Delete{Person} over a "+ + "suspension would destroy an identity the user gets back") + }) + } +} + +func TestAccount_ReactivationUnpauses(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + fixture := newDispatchFixture(t, database) + + require.NoError(t, fixture.handle(t, + accountFrameFor(dispatchNativeDID, false, "deactivated"))) + require.True(t, deliveryPaused(t, database, dispatchNativeDID)) + + require.NoError(t, fixture.handle(t, accountFrameFor(dispatchNativeDID, true, "active"))) + assert.False(t, deliveryPaused(t, database, dispatchNativeDID), + "coming back resumes delivery, under the same identity and the same frozen "+ + "local part") +} + +func TestAccount_DeletedReachesTheTerminalSeamAndNotThePause(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + terminator := &recordingTerminator{} + fixture := newDispatchFixture(t, database, + func(opts *Options) { opts.Terminator = terminator }) + + require.NoError(t, fixture.handle(t, + accountFrameFor(dispatchNativeDID, false, "deleted"))) + + assert.Equal(t, []string{dispatchNativeDID}, terminator.DIDs(), + "status=deleted is the ONE value that means gone, and it goes to the tier that "+ + "re-verifies against PLC and the PDS before sending Delete{Person} — acting "+ + "on a stale deletion event is unrecoverable") + + assert.False(t, deliveryPaused(t, database, dispatchNativeDID), + "and it is NOT a pause: pausing a deleted account would leave the identity "+ + "standing while pretending something had been done about it") +} + +func TestAccount_DeletedWithNoTerminatorWiredChangesNothing(t *testing.T) { + database := dispatchTestDB(t) + seedAPActor(t, database, dispatchNativeDID, "alice") + // A deployment where task 17 has not landed. + fixture := newDispatchFixture(t, database, func(opts *Options) { opts.Terminator = nil }) + + require.NoError(t, fixture.handle(t, + accountFrameFor(dispatchNativeDID, false, "deleted")), + "a missing terminal tier must not panic and must not fail the event") + + assert.False(t, deliveryPaused(t, database, dispatchNativeDID), + "and must not quietly degrade into a pause — a deletion half-handled as a "+ + "pause looks handled in the database and is not") + enabled, _ := actorEnabled(t, database, dispatchNativeDID) + assert.True(t, enabled) +} + +func TestAccount_ActorlessDIDIsSkippedEverywhere(t *testing.T) { + for _, tc := range []struct { + name string + active bool + status string + }{ + {"deactivated", false, "deactivated"}, + {"reactivated", true, "active"}, + {"deleted", false, "deleted"}, + } { + t.Run(tc.name, func(t *testing.T) { + database := dispatchTestDB(t) + terminator := &recordingTerminator{} + fixture := newDispatchFixture(t, database, + func(opts *Options) { opts.Terminator = terminator }) + + require.NoError(t, fixture.handle(t, + accountFrameFor(dispatchNativeDID, tc.active, tc.status)), + "an account event for a DID with no AP identity is a no-op") + + assert.Zero(t, countRows(t, database, "ap_actors"), + "minting an actor to pause or delete it would create the identity the "+ + "event is about losing") + assert.Empty(t, fixture.minter.Handles()) + assert.Empty(t, terminator.DIDs(), + "and there is nothing for the terminal tier to withdraw: nothing was "+ + "ever federated under this DID") + }) + } +} diff --git a/internal/consume/comments.go b/internal/consume/comments.go index d3348a1..4619cff 100644 --- a/internal/consume/comments.go +++ b/internal/consume/comments.go @@ -212,18 +212,9 @@ func (d *Dispatcher) commentThread(ctx context.Context, atURI string, commit *Co return d.resolveParent(ctx, commit) } -// resolveParent finds the thing a new comment replies to. A parent can live in -// EITHER of two places, and both are legitimate: -// -// - ap_objects: content materialized from the fediverse (a Lemmy post or -// comment), or bridge-origin content mapped at write time; -// - outbound_objects: a NATIVE postv2 the acceptance engine admitted, or an -// earlier native comment. Nothing maps those into ap_objects — they were -// never materialized from the fediverse — so their outbound row is the -// only evidence they federate at all. -// -// A nil thread with a nil error means "not federated here": a skip, not a -// failure. +// resolveParent finds the thing a new comment replies to, and places the +// comment one level below it. Reading the parent's RECORDED depth is what +// keeps the cap O(1) instead of walking the thread on every comment. func (d *Dispatcher) resolveParent(ctx context.Context, commit *CommitEvent) (*resolvedThread, error) { parentATURI := replyRef(commit.Record, "parent") if parentATURI == "" { @@ -234,58 +225,19 @@ func (d *Dispatcher) resolveParent(ctx context.Context, commit *CommitEvent) (*r return nil, nil } - mapping, err := d.objectMappings.GetByATURI(ctx, parentATURI) - if err != nil && !errors.IsNotFound(err) { - return nil, fmt.Errorf("resolve comment parent %s: %w", parentATURI, err) - } - if err == nil && mapping.CommunityDID != "" { - community, err := d.communities.GetByDID(ctx, mapping.CommunityDID) - if errors.IsNotFound(err) { - // The parent is mapped but its community is not one this bridge - // federates, so there is nowhere to deliver to. - return nil, nil - } - if err != nil { - return nil, fmt.Errorf("resolve community %s: %w", mapping.CommunityDID, err) - } - return &resolvedThread{ - ParentATURI: parentATURI, - ParentAPID: mapping.APID, - CommunityDID: mapping.CommunityDID, - CommunityAPID: community.APGroupID, - Depth: d.parentDepth(ctx, parentATURI) + 1, - }, nil - } - - state, err := d.objects.GetByATURI(ctx, parentATURI) - if errors.IsNotFound(err) { - return nil, nil - } - if err != nil { - return nil, fmt.Errorf("read parent outbound state %s: %w", parentATURI, err) + parent, err := d.resolveSubject(ctx, parentATURI) + if err != nil || parent == nil { + return nil, err } return &resolvedThread{ - ParentATURI: parentATURI, - ParentAPID: state.APObjectID, - CommunityDID: state.CommunityDID, - CommunityAPID: state.CommunityAPID, - Depth: state.Depth + 1, + ParentATURI: parent.ATURI, + ParentAPID: parent.APID, + CommunityDID: parent.CommunityDID, + CommunityAPID: parent.CommunityAPID, + Depth: parent.Depth + 1, }, nil } -// parentDepth reads a mapped parent's recorded depth, which exists only if the -// bridge federated it outward too. A parent with no outbound row is a post or -// a Lemmy object at the top of what this bridge tracks, so its replies are -// depth 1. Reading the parent's recorded depth is what keeps the cap O(1) -// instead of walking the thread on every comment. -func (d *Dispatcher) parentDepth(ctx context.Context, parentATURI string) int { - state, err := d.objects.GetByATURI(ctx, parentATURI) - if err != nil { - return 0 - } - return state.Depth -} - // ensureActor makes sure the author has an AP identity, resolving their handle // first if they do not. // diff --git a/internal/consume/consume.go b/internal/consume/consume.go index afa3959..5a966ff 100644 --- a/internal/consume/consume.go +++ b/internal/consume/consume.go @@ -14,6 +14,9 @@ package consume import ( "context" "errors" + "fmt" + "net/url" + "strings" ) // ErrPermanentEvent marks a handler failure as permanent: the event can never @@ -45,6 +48,46 @@ const ( CollectionVote = "social.coves.feed.vote" ) +// WantedCollections is the wantedCollections filter this consumer subscribes +// with. It is a function rather than a package var so no caller can append to +// the shared slice and silently widen every future subscription. +func WantedCollections() []string { + return []string{ + CollectionFederation, + CollectionProfile, + CollectionPostV2, + CollectionComment, + CollectionVote, + } +} + +// SubscribeURL builds the WebSocket subscribe URL: the configured base with a +// /subscribe path (appended unless already present) and ONE wantedCollections +// parameter per collection. +// +// The filter is not optional in practice: a subscribe URL without +// wantedCollections asks Jetstream for the entire firehose, which this +// consumer would then discard record by record — the collection filter is the +// difference between a few Coves repos' events and the whole network's. +func SubscribeURL(baseURL string, collections []string) (string, error) { + parsed, err := url.Parse(baseURL) + if err != nil { + return "", fmt.Errorf("invalid Jetstream URL %q: %w", baseURL, err) + } + if !strings.HasSuffix(parsed.Path, subscribePath) { + parsed.Path = strings.TrimSuffix(parsed.Path, "/") + subscribePath + } + query := parsed.Query() + for _, collection := range collections { + query.Add("wantedCollections", collection) + } + parsed.RawQuery = query.Encode() + return parsed.String(), nil +} + +// subscribePath is Jetstream's subscribe endpoint. +const subscribePath = "/subscribe" + // JetstreamEvent is one frame off the Jetstream WebSocket. type JetstreamEvent struct { Account *AccountEvent `json:"account,omitempty"` diff --git a/internal/consume/dispatch.go b/internal/consume/dispatch.go index c2f4cb3..58037fc 100644 --- a/internal/consume/dispatch.go +++ b/internal/consume/dispatch.go @@ -11,6 +11,7 @@ import ( "strings" "tidepool/internal/errors" + "tidepool/internal/materialize" "tidepool/internal/store" ) @@ -64,6 +65,11 @@ type VoteIntent struct { Direction string // ID is the deterministic activity id. ID string + // InnerActivityID is the id of the Like/Dislike an Undo withdraws. An + // Undo has to embed the activity it undoes, and by the time the delete + // arrives the vote record is gone — so this is read back from + // outbound_votes.current_activity_id, not recomputed. + InnerActivityID string // CommunityAPID is the target community's AP Group id. CommunityAPID string } @@ -103,6 +109,18 @@ type RemoteContentDeleter interface { DeleteRemoteContent(ctx context.Context, did string) error } +// AccountTerminator is the task 17 TERMINAL seam: a repo whose account status +// is "deleted" (decision 19) means the user is gone, and their federated +// identity has to be withdrawn with Delete{Person}. +// +// It is a seam rather than an inline write because acting on a stale event is +// unrecoverable: the terminal tier re-verifies against PLC/the PDS before +// sending anything. Optional — a nil terminator means the tier is not wired +// yet, which is announced rather than silently treated as a pause. +type AccountTerminator interface { + TerminateAccount(ctx context.Context, did string) error +} + // Options configures a Dispatcher. type Options struct { // DB is the bridge database. @@ -117,6 +135,13 @@ type Options struct { // RemoteDeleter is the task 17 destructive seam. Optional (see // RemoteContentDeleter). RemoteDeleter RemoteContentDeleter + // Terminator is the task 17 terminal seam for status=deleted accounts. + // Optional (see AccountTerminator). + Terminator AccountTerminator + // Records reads committed records so a subject's community can be derived + // for mappings written before migration 016 filled community_did. + // *repo.Manager satisfies it. Optional — see subjectCommunityDID. + Records materialize.RecordGetter // Resolver verifies a DID's handle before the FIRST mint. Required: the // local part is frozen at creation, so minting without a verified handle // would freeze a guess. @@ -137,6 +162,7 @@ type Dispatcher struct { enqueuer OutboundEnqueuer engine AcceptanceEngine remoteDeleter RemoteContentDeleter + terminator AccountTerminator resolver DIDResolver prefs store.FederationPrefs apActors store.APActors @@ -145,7 +171,9 @@ type Dispatcher struct { // comment's thread and target community are resolved FROM. objectMappings store.APObjects objects store.OutboundObjects + votes store.OutboundVotes communities store.Communities + records materialize.RecordGetter hosted *hostedRepos gate *RevGate userOrigin string @@ -188,12 +216,15 @@ func NewDispatcher(opts Options) (*Dispatcher, error) { enqueuer: opts.Enqueuer, engine: opts.Engine, remoteDeleter: opts.RemoteDeleter, + terminator: opts.Terminator, resolver: opts.Resolver, prefs: store.NewFederationPrefs(opts.DB), apActors: store.NewAPActors(opts.DB), objectMappings: store.NewAPObjects(opts.DB), objects: store.NewOutboundObjects(opts.DB), + votes: store.NewOutboundVotes(opts.DB), communities: store.NewCommunities(opts.DB), + records: opts.Records, hosted: newHostedRepos(opts.DB), gate: NewRevGate(opts.DB), userOrigin: opts.UserOrigin, @@ -246,6 +277,8 @@ func (d *Dispatcher) commitHandlerFor(collection string) commitHandler { return d.handleComment case CollectionProfile: return d.handleProfile + case CollectionVote: + return d.handleVote } return nil } @@ -371,22 +404,44 @@ func (d *Dispatcher) handleAccount(ctx context.Context, event *JetstreamEvent) e did = event.DID } + // 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 + // order to pause or delete it would create the very identity the event is + // about losing. + if _, err := d.apActors.GetByDID(ctx, did); err != nil { + if errors.IsNotFound(err) { + d.logger.Debug("account status for a DID with no actor", + slog.String("did", did), slog.String("status", account.Status)) + return nil + } + return fmt.Errorf("look up actor for %s: %w", did, err) + } + if account.Status == accountStatusDeleted { - // The terminal tier re-verifies against PLC/PDS before sending - // Delete{Person}, because a deletion acted on from a stale event is - // unrecoverable. It is a separate seam and is not wired here. - d.logger.Warn("account reported deleted; the terminal tier owns this event", - slog.String("did", did)) + // The ONE status that means gone. It goes to the tier that re-verifies + // against PLC and the PDS before sending Delete{Person}, because + // acting on a stale deletion event is unrecoverable. + if d.terminator == nil { + // 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 nil + } + if err := d.terminator.TerminateAccount(ctx, did); err != nil { + return fmt.Errorf("terminate account %s: %w", did, err) + } return nil } + // 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. err := d.apActors.SetPaused(ctx, did, !account.Active) if errors.IsNotFound(err) { - // No actor for this DID: there is no delivery to pause, and minting - // one to pause it would create the identity the event is about losing. - d.logger.Debug("account status for a DID with no actor", - slog.String("did", did), slog.String("status", account.Status)) - return nil + return nil // the actor vanished between the check and the write } return err } diff --git a/internal/consume/metrics.go b/internal/consume/metrics.go new file mode 100644 index 0000000..7e4847b --- /dev/null +++ b/internal/consume/metrics.go @@ -0,0 +1,110 @@ +package consume + +import ( + "context" + "expvar" + "sync" + "time" +) + +// Metric names published for /admin/metrics. They carry the "tidepool_" +// prefix because ingest.scopedMetrics serves ONLY that prefix — a gauge named +// without it is published to expvar and then filtered out of the endpoint, +// which looks exactly like a gauge that is not being updated. +const ( + MetricCursorAgeSeconds = "tidepool_consumer_cursor_age_seconds" + MetricLastEventAgeSeconds = "tidepool_consumer_last_event_age_seconds" + MetricDeadLetterDepth = "tidepool_consumer_dead_letters" + MetricReconnects = "tidepool_consumer_reconnects" + MetricEventsProcessed = "tidepool_consumer_events_processed" + MetricEventsDeadLettered = "tidepool_consumer_events_dead_lettered" +) + +// deadLetterDepthUnavailable is what the backlog gauge reports when storage +// cannot be read. A negative value is impossible for a count, so it is +// unmistakable — where a 0 would claim the backlog is empty at exactly the +// moment nobody can tell. +const deadLetterDepthUnavailable = -1 + +// publishOnce guards the expvar registration. expvar PANICS on a duplicate +// name, and both main and the tests call PublishMetrics. +var publishOnce sync.Once + +// PublishMetrics registers the consumer's gauges with expvar. +// +// They are expvar.Func rather than counters the consumer pokes: cursor age and +// dead-letter depth are QUESTIONS about the current state, and a value pushed +// on every event would go stale precisely when the consumer stalls — which is +// the moment an operator looks at it. +// +// A stalled consumer is the failure this task most has to make visible, +// because it is otherwise invisible: the process is up, the health check is +// green, and events simply stop arriving. Cursor age and last-event age are +// the pair that tells a dead upstream from a dead consumer — a quiet stream +// keeps the cursor current while last-event age grows. +// +// Idempotent: a second call is a no-op, so the FIRST connector and queue +// handed in are the ones the gauges read for the life of the process. +func PublishMetrics(ctx context.Context, connector *Connector, queue DeadLetterQueue) { + publishOnce.Do(func() { + expvar.Publish(MetricCursorAgeSeconds, expvar.Func(func() any { + return ageSeconds(cursorTime(connector.Status().CursorTimeUS)) + })) + expvar.Publish(MetricLastEventAgeSeconds, expvar.Func(func() any { + status := connector.Status() + if status.LastEventAt == nil { + return 0.0 + } + return ageSeconds(*status.LastEventAt) + })) + expvar.Publish(MetricDeadLetterDepth, expvar.Func(func() any { + // Read from STORAGE at scrape time, not counted in memory: a + // restart must not reset the backlog to zero and declare it gone. + counts, err := queue.CountDeadLetters(ctx) + if err != nil { + // This gauge is read while somebody is looking at a broken + // system, so a failing storage read reports unavailable rather + // than panicking inside the metrics handler and taking the + // endpoint down with it. + return deadLetterDepthUnavailable + } + var total int64 + for _, count := range counts { + total += count + } + return total + })) + expvar.Publish(MetricReconnects, expvar.Func(func() any { + return connector.Status().Reconnects + })) + expvar.Publish(MetricEventsProcessed, expvar.Func(func() any { + return connector.Status().EventsProcessed + })) + expvar.Publish(MetricEventsDeadLettered, expvar.Func(func() any { + return connector.Status().EventsDeadLettered + })) + }) +} + +// cursorTime converts a Jetstream cursor to wall time. A zero cursor means the +// consumer has never processed an event, which is reported as "no age" rather +// than the decades since the epoch. +func cursorTime(cursorTimeUS int64) time.Time { + if cursorTimeUS <= 0 { + return time.Time{} + } + return time.UnixMicro(cursorTimeUS) +} + +// ageSeconds is how long ago a moment was, floored at zero (clock skew between +// the app and the event source must not produce a negative age). +func ageSeconds(moment time.Time) float64 { + if moment.IsZero() { + return 0 + } + age := time.Since(moment).Seconds() + if age < 0 { + return 0 + } + return age +} diff --git a/internal/consume/metrics_test.go b/internal/consume/metrics_test.go new file mode 100644 index 0000000..86e884a --- /dev/null +++ b/internal/consume/metrics_test.go @@ -0,0 +1,142 @@ +package consume + +import ( + "context" + "expvar" + "strconv" + "strings" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// Task 14 cycle K2: observability. +// +// A stalled consumer is the failure this task most has to make visible, +// because it is otherwise invisible: the process is up, the health check is +// green, and events simply stop arriving. Cursor age and dead-letter depth are +// the two numbers that say so. + +// metricValue reads a published expvar as a float. A missing var fails the +// test rather than returning a zero that would look like a healthy gauge. +func metricValue(t *testing.T, name string) float64 { + t.Helper() + published := expvar.Get(name) + require.NotNil(t, published, "expvar %q must be published", name) + value, err := strconv.ParseFloat(strings.Trim(published.String(), `"`), 64) + require.NoError(t, err, "expvar %q must publish a number, got %s", name, published.String()) + return value +} + +func TestPublishMetrics_ExposesTheGaugesAnOperatorNeeds(t *testing.T) { + database := connectorTestDB(t) + state := NewPostgresStateStore(database, CursorSchemaVersion) + handler := &recordingHandler{} + + // A recent event, so cursor age is a small number rather than the decades + // a synthetic time_us would produce. + nowUS := time.Now().UnixMicro() + fake := newFakeJetstream(t, connFrame(nowUS, "3lzrevmetric1", "aaa")) + connector := newTestConnector(t, fake.URL(), handler, state) + startConnector(t, connector) + + waitFor(t, "the scripted event to be processed", func() bool { return handler.Calls() >= 1 }) + + require.NoError(t, state.AddDeadLetter(context.Background(), ConsumerNative, 1, + []byte(`{"time_us":1}`), "boom", 0)) + + PublishMetrics(context.Background(), connector, state) + + cursorAge := metricValue(t, MetricCursorAgeSeconds) + assert.GreaterOrEqual(t, cursorAge, 0.0) + assert.Less(t, cursorAge, 60.0, + "cursor age is how far behind the stream the consumer is; a consumer that just "+ + "processed a live event must read near zero, or the number is measuring "+ + "something else") + + lastEventAge := metricValue(t, MetricLastEventAgeSeconds) + assert.GreaterOrEqual(t, lastEventAge, 0.0) + assert.Less(t, lastEventAge, 60.0, + "last-event age is the other half: a QUIET stream keeps the cursor current "+ + "while this number grows, which is how a dead upstream is told apart from "+ + "a dead consumer") + + assert.Equal(t, 1.0, metricValue(t, MetricDeadLetterDepth), + "the DLQ backlog is read from storage at scrape time, not counted in memory — "+ + "a restart must not reset it to zero and declare the backlog gone") + assert.Equal(t, 1.0, metricValue(t, MetricEventsProcessed)) + assert.Equal(t, 0.0, metricValue(t, MetricEventsDeadLettered)) + assert.GreaterOrEqual(t, metricValue(t, MetricReconnects), 0.0) +} + +func TestPublishMetrics_NamesAreServedByTheAdminEndpoint(t *testing.T) { + database := connectorTestDB(t) + state := NewPostgresStateStore(database, CursorSchemaVersion) + fake := newFakeJetstream(t) + connector := newTestConnector(t, fake.URL(), &recordingHandler{}, state) + + PublishMetrics(context.Background(), connector, state) + + for _, name := range []string{ + MetricCursorAgeSeconds, MetricLastEventAgeSeconds, MetricDeadLetterDepth, + MetricReconnects, MetricEventsProcessed, MetricEventsDeadLettered, + } { + assert.True(t, strings.HasPrefix(name, "tidepool_"), + "ingest.scopedMetrics serves ONLY the tidepool_ prefix, so a gauge named "+ + "without it is published to expvar and then filtered out of "+ + "/admin/metrics — indistinguishable from a gauge nobody updates. %q", name) + assert.NotNil(t, expvar.Get(name), "%q must be published", name) + } +} + +func TestPublishMetrics_IsIdempotent(t *testing.T) { + database := connectorTestDB(t) + state := NewPostgresStateStore(database, CursorSchemaVersion) + fake := newFakeJetstream(t) + connector := newTestConnector(t, fake.URL(), &recordingHandler{}, state) + + PublishMetrics(context.Background(), connector, state) + assert.NotPanics(t, func() { PublishMetrics(context.Background(), connector, state) }, + "expvar panics on a duplicate name, and both main and the tests call this — a "+ + "second call must be a no-op rather than taking the process down at boot") +} + +func TestPublishMetrics_SurvivesAStoreThatCannotBeReached(t *testing.T) { + database := connectorTestDB(t) + state := NewPostgresStateStore(database, CursorSchemaVersion) + fake := newFakeJetstream(t) + connector := newTestConnector(t, fake.URL(), &recordingHandler{}, state) + + PublishMetrics(context.Background(), connector, failingDeadLetters{}) + + assert.NotPanics(t, func() { _ = expvar.Get(MetricDeadLetterDepth).String() }, + "a gauge is read while an operator is looking at a broken system, so a failing "+ + "storage read must produce a value rather than panic inside the metrics "+ + "handler and take the endpoint down with it") +} + +// failingDeadLetters answers every query with an error — postgres down, which +// is exactly when someone is reading the metrics. +type failingDeadLetters struct{ DeadLetterQueue } + +func (failingDeadLetters) CountDeadLetters(context.Context) (map[string]int64, error) { + return nil, assert.AnError +} + +func TestNewNoopEnqueuer_AcceptsIntentsAndDeliversNothing(t *testing.T) { + enqueuer := NewNoopEnqueuer(nil) + require.NotNil(t, enqueuer, + "a nil logger must default rather than nil-panic on the first intent") + + err := enqueuer.EnqueueActivity(context.Background(), dispatchNativeDID, "key", "at://parent", + CommentIntent{Op: "create", ATURI: "at://x", ID: "https://coves.social/ap/activity/abc"}) + assert.NoError(t, err, + "until task 15 lands, main wires this so the consumer still RUNS and writes its "+ + "durable state — disabling the whole path instead would leave everything "+ + "downstream of the cursor unexercised until delivery exists") + + assert.NoError(t, enqueuer.EnqueueActivity(context.Background(), dispatchNativeDID, "key", "", + VoteIntent{Op: "undo", VoteATURI: "at://v", Direction: "up"})) +} diff --git a/internal/consume/noop.go b/internal/consume/noop.go new file mode 100644 index 0000000..bb1fd71 --- /dev/null +++ b/internal/consume/noop.go @@ -0,0 +1,36 @@ +package consume + +import ( + "context" + "log/slog" +) + +// noopEnqueuer logs the outbound intents it is handed and delivers nothing. +// It is what main wires until task 15 lands, following the v1 precedent +// (ingest.NewNoopVotes): the consumer runs, writes its state, and makes the +// work it WOULD deliver visible, rather than being disabled entirely and +// leaving the whole path unexercised until delivery exists. +type noopEnqueuer struct { + logger *slog.Logger +} + +// NewNoopEnqueuer returns the logging no-op OutboundEnqueuer. A nil logger +// uses slog.Default(). +func NewNoopEnqueuer(logger *slog.Logger) OutboundEnqueuer { + if logger == nil { + logger = slog.Default() + } + return &noopEnqueuer{logger: logger} +} + +// EnqueueActivity records what WOULD have been delivered and drops it. The +// activity id is logged because it is the one field a later delivery has to +// reproduce exactly: it is what a peer dedupes on. +func (e *noopEnqueuer) EnqueueActivity(_ context.Context, actorDID, orderingKey, parentATURI string, intent Intent) error { + e.logger.Info("outbound intent dropped: no delivery queue wired", + slog.String("actor", actorDID), + slog.String("ordering_key", orderingKey), + slog.String("parent", parentATURI), + slog.String("activity", intent.ActivityID())) + return nil +} diff --git a/internal/consume/subjects.go b/internal/consume/subjects.go new file mode 100644 index 0000000..5297bf1 --- /dev/null +++ b/internal/consume/subjects.go @@ -0,0 +1,130 @@ +package consume + +import ( + "context" + "fmt" + + "tidepool/internal/errors" + "tidepool/internal/materialize" + "tidepool/internal/store" +) + +// Subject resolution is the one question a comment and a vote both have to +// answer before anything else: does this bridge federate the thing being +// replied to or voted on, and if so, under which community? +// +// It lives here once because the two paths giving different answers would be a +// fork in what gets delivered where — a comment federated into one community +// while its vote goes to another. + +// resolvedSubject is what the bridge already knows about a thing a native +// record points at. Every field comes from state the bridge itself wrote, +// never from the pointing record. +type resolvedSubject struct { + ATURI string + APID string + // CommunityDID and CommunityAPID are the target community on both sides + // of the bridge. + CommunityDID string + CommunityAPID string + // Depth is the subject's OWN recorded reply depth (0 for a post or a + // subject the bridge does not track depth for). Callers that nest below it + // add one. + Depth int +} + +// resolveSubject looks a subject up in the two places one can live, in order. +// Both are legitimate: +// +// - ap_objects: content materialized FROM the fediverse (a Lemmy post or +// comment), or bridge-origin content mapped at write time; +// - outbound_objects: a NATIVE postv2 the acceptance engine admitted, or an +// earlier native comment. Nothing maps those into ap_objects — they were +// never materialized from the fediverse — so their outbound row is the +// only evidence they federate at all. +// +// A nil subject with a nil error means "not federated here": a skip, not a +// failure. Most native comments and votes land on native content, and +// dead-lettering all of them would bury the queue. +func (d *Dispatcher) resolveSubject(ctx context.Context, atURI string) (*resolvedSubject, error) { + mapping, err := d.objectMappings.GetByATURI(ctx, atURI) + if err != nil && !errors.IsNotFound(err) { + return nil, fmt.Errorf("resolve subject %s: %w", atURI, err) + } + if err == nil { + communityDID, err := d.subjectCommunityDID(ctx, mapping) + if err != nil { + return nil, err + } + if communityDID != "" { + community, err := d.communities.GetByDID(ctx, communityDID) + if errors.IsNotFound(err) { + // Mapped, but its community is not one this bridge federates: + // there is nowhere to deliver to. + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("resolve community %s: %w", communityDID, err) + } + return &resolvedSubject{ + ATURI: atURI, + APID: mapping.APID, + CommunityDID: communityDID, + CommunityAPID: community.APGroupID, + Depth: d.recordedDepth(ctx, atURI), + }, nil + } + } + + state, err := d.objects.GetByATURI(ctx, atURI) + if errors.IsNotFound(err) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("read subject outbound state %s: %w", atURI, err) + } + return &resolvedSubject{ + ATURI: atURI, + APID: state.APObjectID, + CommunityDID: state.CommunityDID, + CommunityAPID: state.CommunityAPID, + Depth: state.Depth, + }, nil +} + +// subjectCommunityDID answers which community a mapped record belongs to. +// +// materialize.CommunityDIDOf is THE answer to that question bridge-wide — +// ingest's announced-delete authorization and votes' announced-vote binding +// both ask it, and a fork between those answers is a fork in who may moderate +// what. It needs a RecordGetter because rows written BEFORE migration 016 have +// no community_did column filled: for those the answer sits in the record +// itself (a postv2's `community` field, a comment's thread root), where no +// UPDATE statement could reach it. +// +// Without a RecordGetter this falls back to the column alone, which is correct +// for everything materialized since 016 and simply blind to the older rows — +// they resolve to "" and their comments and votes are skipped rather than +// misdirected. +func (d *Dispatcher) subjectCommunityDID(ctx context.Context, mapping *store.APObjectMapping) (string, error) { + if d.records == nil { + return mapping.CommunityDID, nil + } + communityDID, err := materialize.CommunityDIDOf(ctx, d.records, mapping) + if err != nil { + return "", fmt.Errorf("resolve community of %s: %w", mapping.ATURI, err) + } + return communityDID, nil +} + +// recordedDepth reads a subject's own reply depth, which exists only if the +// bridge federated it OUTWARD too. A mapped subject with no outbound row is a +// post, or a Lemmy object whose depth this bridge does not track, so it counts +// as the top: replies to it are depth 1. +func (d *Dispatcher) recordedDepth(ctx context.Context, atURI string) int { + state, err := d.objects.GetByATURI(ctx, atURI) + if err != nil { + return 0 + } + return state.Depth +} diff --git a/internal/consume/votes.go b/internal/consume/votes.go new file mode 100644 index 0000000..1f84766 --- /dev/null +++ b/internal/consume/votes.go @@ -0,0 +1,219 @@ +package consume + +import ( + "context" + "fmt" + "log/slog" + + "tidepool/internal/errors" + "tidepool/internal/store" +) + +// The native-vote path: a social.coves.feed.vote record in a Coves user's own +// repo, cast on something this bridge federates. +// +// A vote DELETE commit names the vote record and NOTHING else — not the +// subject, not the direction, not the id the Like went out under. But an Undo +// has to EMBED the activity it withdraws. That gap is the whole reason +// outbound_votes exists, and it is why the row is written before the intent +// and outlives the record it describes. + +// The vote directions this build understands. `direction` is an OPEN enum in +// the lexicon, so an unrecognised value is a forward-compatible record rather +// than a malformed one. +const ( + directionUp = "up" + directionDown = "down" +) + +// handleVote applies one vote commit. +func (d *Dispatcher) handleVote(ctx context.Context, did string, commit *CommitEvent) error { + switch commit.Operation { + case operationCreate, operationUpdate: + return d.applyVoteWrite(ctx, did, commit) + case operationDelete: + return d.applyVoteDelete(ctx, did, commit) + default: + d.logger.Debug("unknown vote operation", + slog.String("operation", commit.Operation), slog.String("did", did)) + return nil + } +} + +// applyVoteWrite records a cast vote and enqueues the Like/Dislike. +// +// The step order mirrors the comment path, and for the same reasons: the +// opt-out gate first (a vote IS a federating interaction, so it mints — an +// earlier draft of this task missed that gate), then everything that decides +// whether the vote can federate at all, and only then the identity and the +// state. +func (d *Dispatcher) applyVoteWrite(ctx context.Context, did string, commit *CommitEvent) error { + federating, err := d.mayFederate(ctx, did) + if err != nil { + return err + } + if !federating { + d.logger.Debug("skipping vote from an opted-out author", + slog.String("did", did), slog.String("rkey", commit.RKey)) + return nil + } + + direction := stringField(commit.Record, "direction") + if direction != directionUp && direction != directionDown { + // Nothing is stored: guessing a direction would push a vote the user + // never cast, and dead-lettering would turn a lexicon rollout into a + // queue full of rows nobody can redrive. + d.logger.Debug("skipping vote with an unrecognised direction", + slog.String("did", did), slog.String("direction", direction)) + return nil + } + + subjectATURI := refURI(commit.Record, "subject") + if subjectATURI == "" { + d.logger.Debug("skipping vote with no subject", slog.String("did", did)) + return nil + } + subject, err := d.resolveSubject(ctx, subjectATURI) + if err != nil { + return err + } + if subject == nil { + // Native users vote in native communities constantly; dead-lettering + // that would bury the queue. + d.logger.Debug("skipping vote on a subject this bridge does not federate", + slog.String("did", did), slog.String("subject", subjectATURI)) + return nil + } + + if err := d.ensureActor(ctx, did); err != nil { + return err + } + + voteATURI := commitRecordURI(did, commit) + // The seq is derived BEFORE the write so the stored id and the intent's id + // are the same string: they must agree, or the Undo would withdraw an + // activity the peer never saw. Upsert bumps from the same base — 0 for a + // new row, +1 for a re-cast — so the two stay in step. + seq := 0 + if existing, err := d.votes.GetByATURI(ctx, voteATURI); err == nil { + seq = existing.ActivitySeq + 1 + } else if !errors.IsNotFound(err) { + return fmt.Errorf("read vote state for %s: %w", voteATURI, err) + } + + stored, err := d.votes.Upsert(ctx, store.OutboundVote{ + VoteATURI: voteATURI, + ActorDID: did, + SubjectATURI: subjectATURI, + SubjectAPID: subject.APID, + CommunityDID: subject.CommunityDID, + Direction: direction, + // Stored, not recomputed at delete time: by then the vote record is + // gone and this id is the only handle on the activity the Undo has to + // name. + CurrentActivityID: ActivityID(d.userOrigin, voteATURI, operationCreate, seq), + DeliveredState: store.DeliveredStatePending, + }) + if errors.IsAlreadyExists(err) { + // A SECOND vote record for a subject this actor already has a live + // vote on. The existing row is left exactly as it is — clobbering it + // would strand the Undo still owed for the Like already delivered. + // + // TRANSIENT, not a skip: the newer record retries until the older + // one's delete lands. A vote change reaching this consumer out of + // order (the delete behind the re-cast) resolves itself on redrive, + // where a skip would drop the user's new vote for good. If the delete + // never comes the budget exhausts and the row is visible in the DLQ, + // which is the right place for two sides disagreeing about what the + // user's vote is. + d.logger.Warn("a second live vote for one subject; retrying until the first is deleted", + slog.String("did", did), + slog.String("vote", voteATURI), + slog.String("subject", subjectATURI)) + return err + } + if err != nil { + return fmt.Errorf("write vote state for %s: %w", voteATURI, err) + } + + return d.enqueuer.EnqueueActivity(ctx, did, did, subjectATURI, VoteIntent{ + Op: operationCreate, + VoteATURI: voteATURI, + SubjectAPID: stored.SubjectAPID, + Direction: stored.Direction, + ID: stored.CurrentActivityID, + CommunityAPID: subject.CommunityAPID, + }) +} + +// applyVoteDelete withdraws a vote, using ONLY state. +// +// Like a comment delete, this is NOT gated on the opt-out: an Undo only ever +// removes something, and blocking it would leave the user's vote standing on +// the peer forever — the opposite of what asking to stop federating means. +func (d *Dispatcher) applyVoteDelete(ctx context.Context, did string, commit *CommitEvent) error { + voteATURI := commitRecordURI(did, commit) + + stored, err := d.votes.GetByATURI(ctx, voteATURI) + if errors.IsNotFound(err) { + // No Undo may be sent for a Like no peer ever received. + d.logger.Debug("skipping delete for a vote with no outbound state", + slog.String("did", did), slog.String("vote", voteATURI)) + return nil + } + if err != nil { + return fmt.Errorf("read vote state for %s: %w", voteATURI, err) + } + + // Re-upserting the row bumps the seq — the Undo is the next activity, and + // its id must not collide with the Like's — while keeping every other + // column, CurrentActivityID above all: that is the id the Like was + // delivered under, and the Undo has to embed it. The row SURVIVES: task 15 + // needs it to retry the Undo and clears it only once delivery succeeds. + undone, err := d.votes.Upsert(ctx, *stored) + if err != nil { + return fmt.Errorf("bump vote state for %s: %w", voteATURI, err) + } + + return d.enqueuer.EnqueueActivity(ctx, did, did, stored.SubjectATURI, VoteIntent{ + Op: operationUndo, + VoteATURI: voteATURI, + SubjectAPID: stored.SubjectAPID, + // Read back from state: the delete commit carries neither, and an + // Undo{Like} withdrawing a Dislike would move the peer's count the + // wrong way. + Direction: stored.Direction, + ID: ActivityID(d.userOrigin, voteATURI, operationUndo, undone.ActivitySeq), + InnerActivityID: stored.CurrentActivityID, + CommunityAPID: d.communityAPID(ctx, stored.CommunityDID), + }) +} + +// 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" + +// communityAPID resolves a community's AP Group id for addressing. A miss +// yields "" rather than an error: the withdrawal still has to go out, and task +// 15 can address it from the subject. +func (d *Dispatcher) communityAPID(ctx context.Context, communityDID string) string { + if communityDID == "" { + return "" + } + community, err := d.communities.GetByDID(ctx, communityDID) + if err != nil { + return "" + } + return community.APGroupID +} + +// refURI reads a strong-ref's uri out of a decoded record. +func refURI(record map[string]any, name string) string { + ref, ok := record[name].(map[string]any) + if !ok { + return "" + } + uri, _ := ref["uri"].(string) + return uri +} diff --git a/internal/consume/votes_test.go b/internal/consume/votes_test.go new file mode 100644 index 0000000..700ee6a --- /dev/null +++ b/internal/consume/votes_test.go @@ -0,0 +1,361 @@ +package consume + +import ( + "context" + "fmt" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "tidepool/internal/store" +) + +// Task 14 cycle I: social.coves.feed.vote. +// +// A vote DELETE commit names the vote record and nothing else — not the +// subject, not the direction, not the id the Like went out under. But an Undo +// has to EMBED the activity it withdraws. That gap is the reason +// outbound_votes exists, and it is why the row is written before the intent +// and outlives the record it describes. + +func voteFrame(did, rev, rkey, subjectATURI, direction string) []byte { + return []byte(fmt.Sprintf( + `{"did":%q,"time_us":9500,"kind":"commit","commit":{"rev":%q,"operation":"create",`+ + `"collection":"social.coves.feed.vote","rkey":%q,`+ + `"cid":"bafyreievgu2ty7qbiaaom5zhmkznsnajuzideek3lo7e65dwqlrvrxnmo4",`+ + `"record":{"$type":"social.coves.feed.vote","subject":{"uri":%q,"cid":%q},`+ + `"direction":%q,"createdAt":"2026-08-13T10:00:00.000Z"}}}`, + did, rev, rkey, subjectATURI, acceptRootCID, direction)) +} + +func voteDeleteFrame(did, rev, rkey string) []byte { + return []byte(fmt.Sprintf( + `{"did":%q,"time_us":9600,"kind":"commit","commit":{"rev":%q,"operation":"delete",`+ + `"collection":"social.coves.feed.vote","rkey":%q}}`, + did, rev, rkey)) +} + +func voteATURIFor(did, rkey string) string { + return "at://" + did + "/" + CollectionVote + "/" + rkey +} + +// --------------------------------------------------------------------------- +// I1 — a vote on a bridged (Lemmy-origin) subject +// --------------------------------------------------------------------------- + +func TestVoteCreate_OnAMappedSubjectWritesStateAndOneIntent(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + fixture := newDispatchFixture(t, database) + ctx := context.Background() + + const rkey = "3lzvote000001" + voteATURI := voteATURIFor(dispatchNativeDID, rkey) + + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "up"))) + + // The author had never federated anything, so this vote earned them an + // identity — through the same verified-handle path a comment uses. + assert.Equal(t, []string{dispatchNativeDID}, fixture.resolver.Calls(), + "a vote is a federating interaction, so it mints — an earlier draft of the "+ + "task missed this gate entirely") + assert.Equal(t, []string{dispatchNativeHandle}, fixture.minter.Handles()) + + stored, err := store.NewOutboundVotes(database).GetByATURI(ctx, voteATURI) + require.NoError(t, err, "the vote's state is keyed by the VOTE record's at-uri, "+ + "because that is the only thing its delete commit will carry") + require.NotNil(t, stored) + + assert.Equal(t, dispatchNativeDID, stored.ActorDID) + assert.Equal(t, acceptRootATURI, stored.SubjectATURI) + assert.Equal(t, acceptRootAPID, stored.SubjectAPID, + "the subject's AP id comes from its ap_objects mapping") + assert.Equal(t, acceptCommunityDID, stored.CommunityDID) + assert.Equal(t, "up", stored.Direction) + assert.Equal(t, 0, stored.ActivitySeq) + assert.Equal(t, store.DeliveredStatePending, stored.DeliveredState, + "the consumer records INTENT only; task 15 claims delivery, and only on success") + assert.Equal(t, ActivityID(acceptUserOrigin, voteATURI, "create", 0), stored.CurrentActivityID, + "the id the Like goes out under is STORED, because the Undo has to embed it "+ + "long after the vote record is gone") + + calls := fixture.enqueuer.Calls() + require.Len(t, calls, 1, "exactly one intent per vote") + assert.Equal(t, dispatchNativeDID, calls[0].ActorDID) + + intent, ok := calls[0].Intent.(VoteIntent) + require.True(t, ok, "want VoteIntent, got %T", calls[0].Intent) + assert.Equal(t, "create", intent.Op) + assert.Equal(t, voteATURI, intent.VoteATURI) + assert.Equal(t, acceptRootAPID, intent.SubjectAPID) + assert.Equal(t, "up", intent.Direction) + assert.Equal(t, acceptCommunityAPID, intent.CommunityAPID) + assert.Equal(t, stored.CurrentActivityID, intent.ActivityID(), + "the intent and the stored state must agree on the id, or the Undo would "+ + "withdraw an activity the peer never saw") +} + +func TestVoteCreate_DownvoteCarriesItsDirection(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + fixture := newDispatchFixture(t, database) + + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, "3lzvote000002", acceptRootATURI, "down"))) + + stored, err := store.NewOutboundVotes(database).GetByATURI( + context.Background(), voteATURIFor(dispatchNativeDID, "3lzvote000002")) + require.NoError(t, err) + require.NotNil(t, stored) + assert.Equal(t, "down", stored.Direction, + "direction is stored verbatim: it is what the Undo is reconstructed from, and "+ + "a Dislike withdrawn as a Like would corrupt the peer's count") +} + +// --------------------------------------------------------------------------- +// I2 — a vote on a native accepted post +// --------------------------------------------------------------------------- + +func TestVoteCreate_OnASubjectKnownOnlyFromOutboundState(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + fixture := newDispatchFixture(t, database) + ctx := context.Background() + + // A native postv2 the acceptance engine admitted: never materialized from + // the fediverse, so no ap_objects mapping exists for it at all. + subjectATURI := "at://" + acceptRootAuthorDID + "/" + CollectionPostV2 + "/3lznativevt1" + subject := seedOutboundParent(t, database, subjectATURI, 0) + require.Zero(t, countRows(t, database, "ap_objects")) + + const rkey = "3lzvote000003" + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, rkey, subjectATURI, "up"))) + + stored, err := store.NewOutboundVotes(database).GetByATURI( + ctx, voteATURIFor(dispatchNativeDID, rkey)) + require.NoError(t, err, + "votes on native accepted posts must federate: those posts ARE the bridged "+ + "content, and dropping their votes would freeze every native thread's score") + require.NotNil(t, stored) + assert.Equal(t, subject.APObjectID, stored.SubjectAPID) + assert.Equal(t, acceptCommunityDID, stored.CommunityDID, + "the community comes from the subject's outbound state — the same answer the "+ + "ap_objects path gives") +} + +func TestVoteCreate_OnAnUnknownSubjectIsSkipped(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + fixture := newDispatchFixture(t, database) + + require.NoError(t, fixture.handle(t, voteFrame(dispatchNativeDID, dispatchRev, + "3lzvote000004", "at://did:plc:nobody/social.coves.community.postv2/nope", "up")), + "a vote on a subject this bridge does not federate is a SKIP: native users vote "+ + "in native communities constantly, and dead-lettering that would bury the queue") + + assert.Zero(t, countRows(t, database, "outbound_votes")) + assert.Empty(t, fixture.enqueuer.Calls()) + assert.Empty(t, fixture.minter.Handles(), "and nothing is minted for it") +} + +// --------------------------------------------------------------------------- +// I3 — the Undo, rebuilt entirely from state +// --------------------------------------------------------------------------- + +func TestVoteDelete_ReconstructsTheUndoFromState(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + fixture := newDispatchFixture(t, database) + ctx := context.Background() + + const rkey = "3lzvote000005" + voteATURI := voteATURIFor(dispatchNativeDID, rkey) + + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "down"))) + created, err := store.NewOutboundVotes(database).GetByATURI(ctx, voteATURI) + require.NoError(t, err) + require.NotNil(t, created) + likeID := created.CurrentActivityID + + // The delete frame carries the vote's at-uri and nothing else. + require.NoError(t, fixture.handle(t, + voteDeleteFrame(dispatchNativeDID, dispatchRevHigher, rkey))) + + calls := fixture.enqueuer.Calls() + require.Len(t, calls, 2, "one Like, one Undo") + intent, ok := calls[1].Intent.(VoteIntent) + require.True(t, ok, "want VoteIntent, got %T", calls[1].Intent) + + assert.Equal(t, "undo", intent.Op) + assert.Equal(t, voteATURI, intent.VoteATURI) + assert.Equal(t, "down", intent.Direction, + "the direction is read back from state — the delete commit does not carry it, "+ + "and an Undo{Like} withdrawing a Dislike would move the peer's count the "+ + "wrong way") + assert.Equal(t, acceptRootAPID, intent.SubjectAPID, "as is the subject") + assert.Equal(t, likeID, intent.InnerActivityID, + "the Undo EMBEDS the activity it withdraws, by the id the Like was delivered "+ + "under — a freshly derived id would name an activity the peer never saw") + + updated, err := store.NewOutboundVotes(database).GetByATURI(ctx, voteATURI) + require.NoError(t, err, "the row SURVIVES the delete: task 15 needs it to retry the "+ + "Undo, and clears it only once delivery succeeds") + require.NotNil(t, updated) + assert.Equal(t, 1, updated.ActivitySeq, "the Undo is the next activity") + assert.Equal(t, ActivityID(acceptUserOrigin, voteATURI, "undo", 1), intent.ActivityID(), + "and its id derives from the bumped seq, so it can never collide with the Like's") + assert.Equal(t, store.DeliveredStatePending, updated.DeliveredState, + "the consumer still claims no delivery; task 15 flips this to undone on success") +} + +func TestVoteDelete_OfAnUnknownVoteIsSkipped(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + fixture := newDispatchFixture(t, database) + + require.NoError(t, fixture.handle(t, + voteDeleteFrame(dispatchNativeDID, dispatchRev, "3lzvote000006")), + "a vote this bridge never federated has nothing to withdraw") + + assert.Empty(t, fixture.enqueuer.Calls(), + "and no Undo may be sent for a Like no peer ever received") +} + +// --------------------------------------------------------------------------- +// I4 — records the handler cannot act on +// --------------------------------------------------------------------------- + +func TestVoteCreate_UnknownDirectionIsSkipped(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + fixture := newDispatchFixture(t, database) + + // knownValues is an OPEN enum: the lexicon may grow a direction this build + // has never heard of. + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, "3lzvote000007", acceptRootATURI, "sideways")), + "RULED SKIP: a direction this build does not understand is a forward-compatible "+ + "record, not a malformed one. Dead-lettering it would turn a lexicon "+ + "rollout into a queue full of rows nobody can redrive") + + assert.Zero(t, countRows(t, database, "outbound_votes"), + "nothing is stored: guessing a direction would push a vote the user never cast") + assert.Empty(t, fixture.enqueuer.Calls()) +} + +func TestVoteCreate_SecondVoteForTheSameSubjectLeavesTheFirstIntact(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + fixture := newDispatchFixture(t, database) + ctx := context.Background() + + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, "3lzvote000008", acceptRootATURI, "up"))) + first, err := store.NewOutboundVotes(database).GetByATURI( + ctx, voteATURIFor(dispatchNativeDID, "3lzvote000008")) + require.NoError(t, err) + require.NotNil(t, first) + + // A SECOND vote record on the same subject while the first still stands. + // One actor holds one live vote per subject, so the store refuses it. + err = fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRevHigher, "3lzvote000009", acceptRootATURI, "down")) + + require.Error(t, err, + "RULED TRANSIENT. Changing a vote is a delete of the old record plus a create "+ + "of the new one, and the create can reach this handler first — on a cursor "+ + "rewind, or because the delete was itself skipped while the author was "+ + "opted out. Skipping the create would lose that vote PERMANENTLY: no event "+ + "ever re-fires it. Retrying costs nothing and succeeds the moment the "+ + "delete lands; a genuinely stuck one exhausts into the DLQ, where it is "+ + "visible instead of silent") + assert.NotErrorIs(t, err, ErrPermanentEvent, + "so the redriver must be allowed to replay it — a permanent classification "+ + "would spend the budget immediately and strand the vote for good") + + assert.Empty(t, storedRev(t, database, voteATURIFor(dispatchNativeDID, "3lzvote000009")), + "and the failed event leaves NO gate row: the redrive replays the same rev, so "+ + "a claimed gate would reject the retry that was supposed to recover it") + + assert.Equal(t, 1, countRows(t, database, "outbound_votes"), + "the second record is not stored") + unchanged, err := store.NewOutboundVotes(database).GetByATURI( + ctx, voteATURIFor(dispatchNativeDID, "3lzvote000008")) + require.NoError(t, err) + require.NotNil(t, unchanged) + assert.Equal(t, "up", unchanged.Direction, + "and the first vote is left exactly as it was — clobbering it would strand the "+ + "Undo that is still owed for the Like already delivered") + assert.Equal(t, first.ActivitySeq, unchanged.ActivitySeq) + + assert.Len(t, fixture.enqueuer.Calls(), 1, "and no second intent goes out") +} + +// --------------------------------------------------------------------------- +// I5 — the opt-out gate, and the retraction asymmetry +// --------------------------------------------------------------------------- + +func TestVoteCreate_IsBlockedForAnOptedOutAuthor(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + ctx := context.Background() + + _, err := store.NewFederationPrefs(database).Upsert(ctx, store.FederationPref{ + DID: dispatchNativeDID, + Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + + fixture := newDispatchFixture(t, database) + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, "3lzvote000010", acceptRootATURI, "up"))) + + assert.Empty(t, fixture.minter.Handles(), + "the opt-out gate runs before the mint here exactly as it does for comments — "+ + "an earlier draft of the task missed this gate on the vote path") + assert.Zero(t, countRows(t, database, "ap_actors")) + assert.Zero(t, countRows(t, database, "outbound_votes")) + assert.Empty(t, fixture.enqueuer.Calls()) +} + +func TestVoteDelete_ProcessesEvenForAnOptedOutAuthor(t *testing.T) { + database := dispatchTestDB(t) + seedBridgedCommunity(t, database) + seedThreadRoot(t, database) + fixture := newDispatchFixture(t, database) + ctx := context.Background() + + const rkey = "3lzvote000011" + require.NoError(t, fixture.handle(t, + voteFrame(dispatchNativeDID, dispatchRev, rkey, acceptRootATURI, "up"))) + require.Len(t, fixture.enqueuer.Calls(), 1) + + // The author opts out AFTER the Like is already out on the fediverse. + _, err := store.NewFederationPrefs(database).Upsert(ctx, store.FederationPref{ + DID: dispatchNativeDID, + Source: store.FederationPrefSourceRecord, + }) + require.NoError(t, err) + + require.NoError(t, fixture.handle(t, + voteDeleteFrame(dispatchNativeDID, dispatchRevHigher, rkey))) + + calls := fixture.enqueuer.Calls() + require.Len(t, calls, 2, + "the same retraction asymmetry as comments: an opt-out stops new content going "+ + "out, but must never trap a Like the user is trying to withdraw — blocking "+ + "the Undo would leave their vote standing on the peer forever") + intent, ok := calls[1].Intent.(VoteIntent) + require.True(t, ok) + assert.Equal(t, "undo", intent.Op) +}