From 79682c882700b56cbeaa4365e73fb0304676cf19 Mon Sep 17 00:00:00 2001 From: Bretton Date: Fri, 10 Jul 2026 23:44:05 -0700 Subject: [PATCH] feat(bridge): consume bridgedStats from bridge-managed repos MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bridged (Lemmy-origin) vote counts now arrive as an optional bridgedStats {upvotes, downvotes, asOf} field on post/comment records emitted by the tidepool bridge. Coves folds them into displayed counts and score so bridged content ranks natively, gated on provenance so only bridge-managed repos may assert aggregate counts. Changes: - lexicons: social.coves.community.post/comment gain #bridgedStats (upvotes, downvotes, asOf — all required within the optional field) - migration 031: bridged_upvote_count / bridged_downvote_count (CHECK >= 0) + bridged_stats_as_of on posts and comments - provenance gate (jetstream/bridge_trust.go): bridgedStats applied only when the repo's stored pds_url matches TRUSTED_BRIDGE_PDS_HOSTS (default-deny; content still indexes with the field ignored otherwise); input hygiene regardless — negatives rejected, 1M magnitude cap mirroring the bridge's MaxSeededCount - post consumer: new UPDATE handler (updates were previously silently dropped — bridged post edits never re-indexed): create-parity security validation, community/author reassignment rejected, soft-delete skip, atomic newer-or-equal asOf guard in SQL, edited_at bumped only on real content changes, RowsAffected checked - comment consumer: same gate + guard on its update path; soft-deleted comments skipped; resurrection recomputes score from surviving native counts - score is inclusive — (native+bridged up) − (native+bridged down) — at every write site; flip-decrement now uses the same clamped arithmetic as the count columns; hot/top sort expressions, cursors, and indexes untouched - read paths fold bridged counts into displayed up/down everywhere (incl. comment GetByURI; the consumer's guard reads raw columns) - log hygiene: infra errors no longer logged as security rejections; false "Jetstream will replay" comments corrected (connectors are log-and-drop, no cursor) Deploy note: TRUSTED_BRIDGE_PDS_HOSTS must name the bridge's PDS host, or bridgedStats is ignored everywhere (safe default). Co-Authored-By: Claude Fable 5 --- cmd/server/main.go | 27 +- internal/atproto/jetstream/bridge_trust.go | 116 +++ .../atproto/jetstream/bridged_stats_test.go | 786 ++++++++++++++++++ .../atproto/jetstream/comment_consumer.go | 215 ++++- internal/atproto/jetstream/post_consumer.go | 371 +++++++-- internal/atproto/jetstream/vote_consumer.go | 24 +- .../social/coves/community/comment.json | 27 + .../lexicon/social/coves/community/post.json | 27 + internal/core/comments/comment.go | 21 +- internal/core/posts/post.go | 7 + .../migrations/031_add_bridged_vote_stats.sql | 57 ++ internal/db/postgres/comment_repo.go | 67 +- internal/db/postgres/discover_repo.go | 4 +- internal/db/postgres/feed_repo.go | 4 +- internal/db/postgres/post_repo.go | 9 +- internal/db/postgres/timeline_repo.go | 4 +- 16 files changed, 1637 insertions(+), 129 deletions(-) create mode 100644 internal/atproto/jetstream/bridge_trust.go create mode 100644 internal/atproto/jetstream/bridged_stats_test.go create mode 100644 internal/db/migrations/031_add_bridged_vote_stats.sql diff --git a/cmd/server/main.go b/cmd/server/main.go index d36cff8..e9d4500 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -749,7 +749,29 @@ func main() { postJetstreamURL = "ws://localhost:6008/subscribe?wantedCollections=social.coves.community.post" } - postEventConsumer := jetstream.NewPostEventConsumer(postRepo, communityRepo, userService, db) + // Provenance gate for bridge-asserted vote aggregates (bridgedStats). Only records + // whose repo is hosted on a trusted bridge PDS may inflate their displayed vote + // counts/score; every other (native) repo is default-denied so it cannot self-assert + // bridgedStats. Configured via TRUSTED_BRIDGE_PDS_HOSTS (comma-separated PDS host + // URLs), mirroring the COMMUNITY_CREATORS allowlist convention. Empty => bridgedStats + // are universally ignored (safe default for deployments with no bridge). + var trustedBridgePDSHosts []string + if hosts := os.Getenv("TRUSTED_BRIDGE_PDS_HOSTS"); hosts != "" { + for _, h := range strings.Split(hosts, ",") { + if h = strings.TrimSpace(h); h != "" { + trustedBridgePDSHosts = append(trustedBridgePDSHosts, h) + } + } + } + bridgeTrust := jetstream.NewBridgeTrust(trustedBridgePDSHosts) + if len(trustedBridgePDSHosts) > 0 { + log.Printf("bridgedStats provenance: trusting %d bridge PDS host(s)", len(trustedBridgePDSHosts)) + } else { + log.Println("bridgedStats provenance: no trusted bridge PDS hosts configured; bridgedStats will be ignored") + } + + postEventConsumer := jetstream.NewPostEventConsumer(postRepo, communityRepo, userService, db, + jetstream.WithPostBridgeTrust(bridgeTrust)) postJetstreamConnector := jetstream.NewPostJetstreamConnector(postEventConsumer, postJetstreamURL) go func() { @@ -817,7 +839,8 @@ func main() { commentJetstreamURL = "ws://localhost:6008/subscribe?wantedCollections=social.coves.community.comment" } - commentEventConsumer := jetstream.NewCommentEventConsumer(commentRepo, db) + commentEventConsumer := jetstream.NewCommentEventConsumer(commentRepo, db, + jetstream.WithCommentBridgeTrust(bridgeTrust)) commentJetstreamConnector := jetstream.NewCommentJetstreamConnector(commentEventConsumer, commentJetstreamURL) go func() { diff --git a/internal/atproto/jetstream/bridge_trust.go b/internal/atproto/jetstream/bridge_trust.go new file mode 100644 index 0000000..ec166f9 --- /dev/null +++ b/internal/atproto/jetstream/bridge_trust.go @@ -0,0 +1,116 @@ +package jetstream + +import ( + "log" + "net/url" + "strings" + "time" +) + +// maxBridgedCount mirrors the tidepool bridge's own MaxSeededCount ceiling +// (1,000,000). A bridgedStats aggregate asserting a count above this is treated as +// malformed/hostile and the entire aggregate is ignored — it is far larger than any +// plausible origin-platform post score and is the shape a score-inflation attack +// takes. Keep this in sync with the bridge if it ever raises MaxSeededCount. +const maxBridgedCount = 1_000_000 + +// BridgeTrust is the provenance gate for bridge-asserted vote aggregates +// (bridgedStats). bridgedStats let a record declare origin-platform vote counts that +// the consumers fold into the displayed counts, denormalized score, and hot-rank. +// Because posts live in community repos and comments live in user repos, ANY native +// repo could otherwise self-assert bridgedStats{upvotes: 10^9} and inflate its own +// ranking at zero cost. BridgeTrust default-denies that: bridgedStats are honoured +// only when the repo that carried them is hosted on a configured trusted bridge PDS. +// +// Why PDS-host allowlist (and not a DID allowlist or live resolution): +// - Coves already resolves and persists each repo's PDS host at index time +// (users.pds_url for a commenter, communities.pds_url for a post's community), +// derived from identity resolution. The gate reuses that stored value, so it needs +// NO new live PLC/identity lookups on the hot Jetstream path. +// - A host allowlist is O(number of bridge instances) config, not O(number of +// bridged actors) — the bridge mints a new DID per federated community/user, so a +// DID allowlist would be unmaintainable. Trusting the bridge's PDS host(s) trusts +// exactly the infrastructure the operator controls. +// - The set is operator config (TRUSTED_BRIDGE_PDS_HOSTS), mirroring the existing +// COMMUNITY_CREATORS allowlist convention (comma-separated env var). +// +// Trade-off / trust assumption (documented deliberately): this trusts that the stored +// pds_url faithfully reflects where the repo is hosted. A hostile actor who could get +// their repo hosted on (or spoof) the bridge PDS host would be trusted — but that +// requires compromising the operator's own bridge infrastructure, which is outside the +// self-assertion threat this gate closes. An empty allowlist (no bridge configured) +// means bridgedStats are universally ignored, which is the safe default. +type BridgeTrust struct { + // hosts holds normalized (scheme+host, lowercased) trusted PDS host keys. + hosts map[string]struct{} +} + +// NewBridgeTrust builds a provenance gate from a list of trusted bridge PDS host URLs +// (typically parsed from the TRUSTED_BRIDGE_PDS_HOSTS env var). Blank/garbage entries +// are dropped. A gate with no usable hosts trusts nothing (default-deny). +func NewBridgeTrust(pdsHosts []string) *BridgeTrust { + m := make(map[string]struct{}, len(pdsHosts)) + for _, h := range pdsHosts { + if n := normalizePDSHost(h); n != "" { + m[n] = struct{}{} + } + } + return &BridgeTrust{hosts: m} +} + +// TrustsPDS reports whether pdsURL is a configured trusted bridge PDS host. +// Default-deny: a nil gate, an empty allowlist, or an empty/unparseable pdsURL all +// return false so bridgedStats are ignored unless provenance is affirmatively proven. +func (b *BridgeTrust) TrustsPDS(pdsURL string) bool { + if b == nil || len(b.hosts) == 0 { + return false + } + n := normalizePDSHost(pdsURL) + if n == "" { + return false + } + _, ok := b.hosts[n] + return ok +} + +// normalizePDSHost reduces a PDS URL to a stable scheme+host comparison key so that +// "https://Bridge.Example/", "https://bridge.example" and "https://bridge.example:443" +// compare consistently. Values that do not parse as a URL with a host fall back to the +// trimmed, lowercased, trailing-slash-stripped string so a plain host in config still +// matches an identically-formatted stored value. +func normalizePDSHost(raw string) string { + s := strings.TrimSpace(raw) + if s == "" { + return "" + } + if u, err := url.Parse(s); err == nil && u.Host != "" { + return strings.ToLower(u.Scheme + "://" + u.Host) + } + return strings.ToLower(strings.TrimRight(s, "/")) +} + +// validatedBridgedStats applies input hygiene to a bridgedStats aggregate and parses +// its asOf timestamp. It returns ok=false — meaning the caller must ignore the WHOLE +// aggregate — when either count is negative, either count exceeds maxBridgedCount, or +// asOf is unparseable. Treating asOf as part of the atomic trio (fix for the create +// path) ensures counts are never applied with a NULL/unknown asOf, which would defeat +// the regression guard. This hygiene runs regardless of provenance; the trust gate is +// a separate, prior check the caller performs first. +func validatedBridgedStats(stats *BridgedStatsFromJetstream, uri string) (up, down int, asOf time.Time, ok bool) { + if stats.Upvotes < 0 || stats.Downvotes < 0 { + log.Printf("Warning: ignoring bridgedStats with negative counts for %s (up=%d down=%d)", + uri, stats.Upvotes, stats.Downvotes) + return 0, 0, time.Time{}, false + } + if stats.Upvotes > maxBridgedCount || stats.Downvotes > maxBridgedCount { + log.Printf("Warning: ignoring bridgedStats exceeding cap %d for %s (up=%d down=%d)", + maxBridgedCount, uri, stats.Upvotes, stats.Downvotes) + return 0, 0, time.Time{}, false + } + t, err := parseBridgedAsOf(stats.AsOf, uri) + if err != nil { + // parseBridgedAsOf already logs the parse failure. + return 0, 0, time.Time{}, false + } + return stats.Upvotes, stats.Downvotes, t, true +} diff --git a/internal/atproto/jetstream/bridged_stats_test.go b/internal/atproto/jetstream/bridged_stats_test.go new file mode 100644 index 0000000..4b7b15f --- /dev/null +++ b/internal/atproto/jetstream/bridged_stats_test.go @@ -0,0 +1,786 @@ +package jetstream + +import ( + "context" + "database/sql" + "os" + "testing" + "time" + + "Coves/internal/core/users" + "Coves/internal/db/postgres" + + _ "github.com/lib/pq" + "github.com/pressly/goose/v3" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// These tests exercise the bridged-vote-stats support end-to-end against a local +// Postgres test database (the same container the rest of the postgres package tests +// use, port 5434 by default / TEST_DATABASE_URL). They are strictly local-only: no +// public PLC/relay/PDS/image hosts are contacted. + +const ( + bridgedTestPrefix = "did:plc:brtest" + bridgedTestCommunity = "did:plc:brtestcommunity" + bridgedTestAuthor = "did:plc:brtestauthor" + bridgedTestOther = "did:plc:brtestotherauthor" + bridgedTestVoter = "did:plc:brtestvoter" + bridgedTestCommenter = bridgedTestPrefix + "commenter" + + // bridgedTestPDS is the trusted bridge PDS host used across these tests. Test + // users/communities are created with this pds_url and the consumers are constructed + // trusting it, so the provenance gate lets their bridgedStats through. Tests that + // exercise the default-deny path override the repo's pds_url instead. + bridgedTestPDS = "https://bridge.test" + bridgedTestNativePDS = "https://native.pds.test" + + asOfEarly = "2026-01-01T00:00:00Z" + asOfLate = "2026-06-01T00:00:00Z" +) + +// bridgeTrustForTests trusts only the bridge PDS host, so records from repos hosted +// there may assert bridgedStats while every other repo is default-denied. +func bridgeTrustForTests() *BridgeTrust { + return NewBridgeTrust([]string{bridgedTestPDS}) +} + +// setupBridgedTestDB connects to the local test database and runs migrations. +func setupBridgedTestDB(t *testing.T) *sql.DB { + t.Helper() + dsn := os.Getenv("TEST_DATABASE_URL") + if dsn == "" { + dsn = "postgres://test_user:test_password@localhost:5434/coves_test?sslmode=disable" + } + db, err := sql.Open("postgres", dsn) + require.NoError(t, err, "Failed to connect to test database") + if pingErr := db.Ping(); pingErr != nil { + _ = db.Close() + t.Skipf("test database not reachable (%v); start it with `make test-db-reset`", pingErr) + } + require.NoError(t, goose.Up(db, "../../db/migrations"), "Failed to run migrations") + return db +} + +func cleanupBridgedTestData(t *testing.T, db *sql.DB) { + t.Helper() + _, _ = db.Exec("DELETE FROM votes WHERE voter_did LIKE $1", bridgedTestPrefix+"%") + _, _ = db.Exec("DELETE FROM comments WHERE commenter_did LIKE $1 OR root_uri LIKE $2", bridgedTestPrefix+"%", "at://"+bridgedTestPrefix+"%") + _, _ = db.Exec("DELETE FROM posts WHERE community_did LIKE $1", bridgedTestPrefix+"%") + _, _ = db.Exec("DELETE FROM communities WHERE did LIKE $1", bridgedTestPrefix+"%") + _, _ = db.Exec("DELETE FROM users WHERE did LIKE $1", bridgedTestPrefix+"%") +} + +// insertBridgedUser inserts a user hosted on the trusted bridge PDS. +func insertBridgedUser(t *testing.T, db *sql.DB, did, handle string) { + t.Helper() + insertBridgedUserOnPDS(t, db, did, handle, bridgedTestPDS) +} + +// insertBridgedUserOnPDS inserts a user with an explicit pds_url, used to exercise the +// provenance gate (a non-bridge pds_url must cause bridgedStats to be ignored). +func insertBridgedUserOnPDS(t *testing.T, db *sql.DB, did, handle, pdsURL string) { + t.Helper() + _, err := db.Exec(`INSERT INTO users (did, handle, pds_url, created_at) VALUES ($1, $2, $3, NOW()) + ON CONFLICT (did) DO UPDATE SET pds_url = EXCLUDED.pds_url`, + did, handle, pdsURL) + require.NoError(t, err) +} + +func insertBridgedCommunity(t *testing.T, db *sql.DB, did, handle, ownerDID string) { + t.Helper() + insertBridgedCommunityOnPDS(t, db, did, handle, ownerDID, bridgedTestPDS) +} + +func insertBridgedCommunityOnPDS(t *testing.T, db *sql.DB, did, handle, ownerDID, pdsURL string) { + t.Helper() + _, err := db.Exec(`INSERT INTO communities (did, handle, name, owner_did, created_by_did, hosted_by_did, pds_url, created_at) + VALUES ($1, $2, $3, $4, $4, $4, $5, NOW()) ON CONFLICT (did) DO UPDATE SET pds_url = EXCLUDED.pds_url`, + did, handle, "Bridged Test Community", ownerDID, pdsURL) + require.NoError(t, err) +} + +func newPostConsumer(t *testing.T, db *sql.DB) *PostEventConsumer { + t.Helper() + postRepo := postgres.NewPostRepository(db) + communityRepo := postgres.NewCommunityRepository(db) + us := newMockUserService() + us.users[bridgedTestAuthor] = &users.User{DID: bridgedTestAuthor, Handle: "brauthor.test", PDSURL: bridgedTestPDS} + // Register the alternate author so update-time validation passes and the + // reassignment-rejection branch (not the author-not-found branch) is exercised. + us.users[bridgedTestOther] = &users.User{DID: bridgedTestOther, Handle: "brother.test", PDSURL: bridgedTestPDS} + return NewPostEventConsumer(postRepo, communityRepo, us, db, WithPostBridgeTrust(bridgeTrustForTests())) +} + +func newCommentConsumer(db *sql.DB) *CommentEventConsumer { + return NewCommentEventConsumer(postgres.NewCommentRepository(db), db, WithCommentBridgeTrust(bridgeTrustForTests())) +} + +func newVoteConsumer(db *sql.DB) *VoteEventConsumer { + return NewVoteEventConsumer(postgres.NewVoteRepository(db), newMockUserService(), db) +} + +func bridgedStatsRecord(up, down int, asOf string) map[string]interface{} { + return map[string]interface{}{ + "upvotes": up, + "downvotes": down, + "asOf": asOf, + } +} + +func postCommitEvent(op, rkey, cid string, record map[string]interface{}) *JetstreamEvent { + return &JetstreamEvent{ + Kind: "commit", + Did: bridgedTestCommunity, + Commit: &CommitEvent{ + Operation: op, + Collection: "social.coves.community.post", + RKey: rkey, + CID: cid, + Record: record, + }, + } +} + +func postRecord(title, content string, bridged map[string]interface{}) map[string]interface{} { + rec := map[string]interface{}{ + "$type": "social.coves.community.post", + "community": bridgedTestCommunity, + "author": bridgedTestAuthor, + "title": title, + "content": content, + "createdAt": "2026-01-01T00:00:00Z", + } + if bridged != nil { + rec["bridgedStats"] = bridged + } + return rec +} + +// readPostRow returns the stored native/bridged columns for assertions. +func readPostRow(t *testing.T, db *sql.DB, uri string) (up, down, bridgedUp, bridgedDown, score int, asOf *time.Time, deletedAt *time.Time, title string) { + t.Helper() + err := db.QueryRow(`SELECT upvote_count, downvote_count, bridged_upvote_count, bridged_downvote_count, score, bridged_stats_as_of, deleted_at, COALESCE(title,'') + FROM posts WHERE uri = $1`, uri). + Scan(&up, &down, &bridgedUp, &bridgedDown, &score, &asOf, &deletedAt, &title) + require.NoError(t, err) + return +} + +func TestPostConsumer_Create_WithBridgedStats(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/create1" + + err := c.HandleEvent(ctx, postCommitEvent("create", "create1", "bafcreate1", + postRecord("Hello", "world", bridgedStatsRecord(10, 3, asOfEarly)))) + require.NoError(t, err) + + up, down, bUp, bDown, score, asOf, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 0, up, "native upvotes untouched") + assert.Equal(t, 0, down, "native downvotes untouched") + assert.Equal(t, 10, bUp) + assert.Equal(t, 3, bDown) + assert.Equal(t, 7, score, "score = (0+10)-(0+3)") + require.NotNil(t, asOf) + + // Read-path fold: displayed stats include bridged counts. + repo := postgres.NewPostRepository(db) + views, err := repo.GetViewsByURIs(ctx, []string{uri}) + require.NoError(t, err) + view := views[uri] + require.NotNil(t, view) + assert.Equal(t, 10, view.Stats.Upvotes, "displayed upvotes folded") + assert.Equal(t, 3, view.Stats.Downvotes, "displayed downvotes folded") + assert.Equal(t, 7, view.Stats.Score) +} + +func TestPostConsumer_Update_BridgedStats_NewerAsOfApplied(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/upd1" + + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "upd1", "bafupd1", + postRecord("t", "c", bridgedStatsRecord(5, 1, asOfEarly))))) + + // Newer asOf -> applied, content + edited_at updated. + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("update", "upd1", "bafupd1b", + postRecord("t2", "c2", bridgedStatsRecord(20, 4, asOfLate))))) + + _, _, bUp, bDown, score, _, _, title := readPostRow(t, db, uri) + assert.Equal(t, 20, bUp) + assert.Equal(t, 4, bDown) + assert.Equal(t, 16, score, "score = (0+20)-(0+4)") + assert.Equal(t, "t2", title, "content updated on edit") + + var editedAt *time.Time + require.NoError(t, db.QueryRow(`SELECT edited_at FROM posts WHERE uri=$1`, uri).Scan(&editedAt)) + assert.NotNil(t, editedAt, "edited_at set on update") +} + +func TestPostConsumer_Update_StrictlyOlderIgnored_EqualApplied(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/upd2" + + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "upd2", "bafupd2", + postRecord("t", "c", bridgedStatsRecord(20, 4, asOfLate))))) + + // Strictly-older asOf -> ignored (bridged counts unchanged). + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("update", "upd2", "bafupd2b", + postRecord("t", "c", bridgedStatsRecord(1, 1, asOfEarly))))) + _, _, bUp, bDown, score, _, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 20, bUp, "strictly-older asOf must not overwrite") + assert.Equal(t, 4, bDown) + assert.Equal(t, 16, score) + + // Equal asOf with DIFFERENT counts -> APPLIED (newer-or-equal guard). tidepool + // truncates asOf, so equal-string collisions carrying fresher counts must not be + // dropped; genuine replays carry identical counts and remain idempotent. + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("update", "upd2", "bafupd2c", + postRecord("t", "c", bridgedStatsRecord(99, 5, asOfLate))))) + _, _, bUp2, bDown2, score2, _, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 99, bUp2, "equal asOf with fresher counts must apply") + assert.Equal(t, 5, bDown2) + assert.Equal(t, 94, score2, "score = (0+99)-(0+5)") +} + +func TestPostConsumer_Update_ReassignmentRejected(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedUser(t, db, bridgedTestOther, "brother.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/reassign1" + + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "reassign1", "bafr1", + postRecord("orig", "c", nil)))) + + // Update attempts to reassign author -> skipped (no error), row unchanged. + rec := postRecord("changed", "c", nil) + rec["author"] = bridgedTestOther + err := c.HandleEvent(ctx, postCommitEvent("update", "reassign1", "bafr1b", rec)) + require.NoError(t, err, "reassignment is skipped, not errored") + + var storedAuthor, storedTitle string + require.NoError(t, db.QueryRow(`SELECT author_did, COALESCE(title,'') FROM posts WHERE uri=$1`, uri).Scan(&storedAuthor, &storedTitle)) + assert.Equal(t, bridgedTestAuthor, storedAuthor, "author must not change") + assert.Equal(t, "orig", storedTitle, "content must not change when reassignment rejected") +} + +func TestPostConsumer_Update_SoftDeletedSkipped(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/del1" + + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "del1", "bafd1", + postRecord("orig", "c", bridgedStatsRecord(5, 0, asOfEarly))))) + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("delete", "del1", "", nil))) + + // Update on soft-deleted row -> skipped, stays deleted, content unchanged. + err := c.HandleEvent(ctx, postCommitEvent("update", "del1", "bafd1b", + postRecord("changed", "c2", bridgedStatsRecord(99, 9, asOfLate)))) + require.NoError(t, err) + + _, _, bUp, _, _, _, deletedAt, title := readPostRow(t, db, uri) + assert.NotNil(t, deletedAt, "post stays soft-deleted") + assert.Equal(t, 5, bUp, "bridged counts untouched on soft-deleted row") + assert.Equal(t, "orig", title, "content untouched on soft-deleted row") +} + +func TestPostConsumer_Update_NonExistentSkipped(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/ghost1" + + err := c.HandleEvent(ctx, postCommitEvent("update", "ghost1", "bafg1", + postRecord("t", "c", bridgedStatsRecord(1, 1, asOfLate)))) + require.NoError(t, err, "update for non-indexed post is a logged skip") + + var count int + require.NoError(t, db.QueryRow(`SELECT COUNT(*) FROM posts WHERE uri=$1`, uri).Scan(&count)) + assert.Equal(t, 0, count, "no row created by update") +} + +func TestPostConsumer_InclusiveScore_NativeVotesStackOnBridged(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + pc := newPostConsumer(t, db) + vc := newVoteConsumer(db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/score1" + + require.NoError(t, pc.HandleEvent(ctx, postCommitEvent("create", "score1", "bafs1", + postRecord("t", "c", bridgedStatsRecord(10, 2, asOfEarly))))) + + // A native upvote stacks on the bridged counts. + voteEvent := &JetstreamEvent{ + Kind: "commit", + Did: bridgedTestVoter, + Commit: &CommitEvent{ + Operation: "create", + Collection: "social.coves.feed.vote", + RKey: "v1", + CID: "bafvote1", + Record: map[string]interface{}{ + "subject": map[string]interface{}{"uri": uri, "cid": "bafs1"}, + "direction": "up", + "createdAt": "2026-02-01T00:00:00Z", + }, + }, + } + require.NoError(t, vc.HandleEvent(ctx, voteEvent)) + + up, down, bUp, bDown, score, _, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 1, up) + assert.Equal(t, 0, down) + assert.Equal(t, 10, bUp) + assert.Equal(t, 2, bDown) + assert.Equal(t, 9, score, "score = (1+10)-(0+2)") + + // A newer bridgedStats update must preserve the native vote. + require.NoError(t, pc.HandleEvent(ctx, postCommitEvent("update", "score1", "bafs1b", + postRecord("t", "c", bridgedStatsRecord(30, 2, asOfLate))))) + up2, _, bUp2, _, score2, _, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 1, up2, "native upvote preserved across bridged update") + assert.Equal(t, 30, bUp2) + assert.Equal(t, 29, score2, "score = (1+30)-(0+2)") + + // Displayed stats fold native + bridged. + repo := postgres.NewPostRepository(db) + views, err := repo.GetViewsByURIs(ctx, []string{uri}) + require.NoError(t, err) + assert.Equal(t, 31, views[uri].Stats.Upvotes) + assert.Equal(t, 2, views[uri].Stats.Downvotes) + assert.Equal(t, 29, views[uri].Stats.Score) +} + +// --- Comment consumer --- + +func commentRecord(content, rootURI, rootCID, parentURI, parentCID string, bridged map[string]interface{}) map[string]interface{} { + rec := map[string]interface{}{ + "$type": "social.coves.community.comment", + "content": content, + "reply": map[string]interface{}{ + "root": map[string]interface{}{"uri": rootURI, "cid": rootCID}, + "parent": map[string]interface{}{"uri": parentURI, "cid": parentCID}, + }, + "createdAt": "2026-01-02T00:00:00Z", + } + if bridged != nil { + rec["bridgedStats"] = bridged + } + return rec +} + +func commentCommitEvent(op, rkey, cid string, record map[string]interface{}) *JetstreamEvent { + return &JetstreamEvent{ + Kind: "commit", + Did: bridgedTestPrefix + "commenter", + Commit: &CommitEvent{ + Operation: op, + Collection: CommentCollection, + RKey: rkey, + CID: cid, + Record: record, + }, + } +} + +func readCommentRow(t *testing.T, db *sql.DB, uri string) (up, down, bridgedUp, bridgedDown, score int, asOf *time.Time) { + t.Helper() + err := db.QueryRow(`SELECT upvote_count, downvote_count, bridged_upvote_count, bridged_downvote_count, score, bridged_stats_as_of + FROM comments WHERE uri=$1`, uri).Scan(&up, &down, &bridgedUp, &bridgedDown, &score, &asOf) + require.NoError(t, err) + return +} + +// setupCommentThread creates a post (so the comment has a valid parent) and returns its URI/CID. +// It also indexes the commenter as a user hosted on the trusted bridge PDS, so the +// comment provenance gate (which resolves the commenter's pds_url) admits bridgedStats. +func setupCommentThread(t *testing.T, db *sql.DB) (postURI, postCID string) { + t.Helper() + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedUser(t, db, bridgedTestCommenter, "brcommenter.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + pc := newPostConsumer(t, db) + postURI = "at://" + bridgedTestCommunity + "/social.coves.community.post/cthread" + postCID = "bafcthread" + require.NoError(t, pc.HandleEvent(context.Background(), + postCommitEvent("create", "cthread", postCID, postRecord("t", "c", nil)))) + return postURI, postCID +} + +func TestCommentConsumer_Create_WithBridgedStats(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + postURI, postCID := setupCommentThread(t, db) + cc := newCommentConsumer(db) + ctx := context.Background() + uri := "at://" + bridgedTestPrefix + "commenter/social.coves.community.comment/cc1" + + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("create", "cc1", "bafcc1", + commentRecord("hi", postURI, postCID, postURI, postCID, bridgedStatsRecord(5, 1, asOfEarly))))) + + up, down, bUp, bDown, score, asOf := readCommentRow(t, db, uri) + assert.Equal(t, 0, up) + assert.Equal(t, 0, down) + assert.Equal(t, 5, bUp) + assert.Equal(t, 1, bDown) + assert.Equal(t, 4, score) + require.NotNil(t, asOf) + + // Read-path fold via ListByParent. + repo := postgres.NewCommentRepository(db) + list, err := repo.ListByParent(ctx, postURI, 10, 0) + require.NoError(t, err) + require.Len(t, list, 1) + assert.Equal(t, 5, list[0].UpvoteCount, "displayed upvotes folded") + assert.Equal(t, 1, list[0].DownvoteCount, "displayed downvotes folded") + assert.Equal(t, 4, list[0].Score) +} + +func TestCommentConsumer_Update_AsOfGuard_AndInclusiveScore(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + postURI, postCID := setupCommentThread(t, db) + cc := newCommentConsumer(db) + vc := newVoteConsumer(db) + ctx := context.Background() + uri := "at://" + bridgedTestPrefix + "commenter/social.coves.community.comment/cc2" + + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("create", "cc2", "bafcc2", + commentRecord("hi", postURI, postCID, postURI, postCID, bridgedStatsRecord(5, 1, asOfEarly))))) + + // Native upvote stacks on bridged. + require.NoError(t, vc.HandleEvent(ctx, &JetstreamEvent{ + Kind: "commit", Did: bridgedTestVoter, + Commit: &CommitEvent{Operation: "create", Collection: "social.coves.feed.vote", RKey: "cv1", CID: "bafcv1", + Record: map[string]interface{}{ + "subject": map[string]interface{}{"uri": uri, "cid": "bafcc2"}, + "direction": "up", "createdAt": "2026-02-01T00:00:00Z", + }}, + })) + up, _, _, _, score, _ := readCommentRow(t, db, uri) + assert.Equal(t, 1, up) + assert.Equal(t, 5, score, "score = (1+5)-(0+1)") + + // Stale asOf update -> ignored, native vote preserved. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("update", "cc2", "bafcc2b", + commentRecord("edited", postURI, postCID, postURI, postCID, bridgedStatsRecord(0, 0, "2025-01-01T00:00:00Z"))))) + up2, _, bUp2, bDown2, score2, _ := readCommentRow(t, db, uri) + assert.Equal(t, 1, up2, "native vote preserved") + assert.Equal(t, 5, bUp2, "stale bridged ignored") + assert.Equal(t, 1, bDown2) + assert.Equal(t, 5, score2) + + // Newer asOf update -> applied, native vote still preserved. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("update", "cc2", "bafcc2c", + commentRecord("edited2", postURI, postCID, postURI, postCID, bridgedStatsRecord(12, 3, asOfLate))))) + up3, _, bUp3, bDown3, score3, _ := readCommentRow(t, db, uri) + assert.Equal(t, 1, up3) + assert.Equal(t, 12, bUp3) + assert.Equal(t, 3, bDown3) + assert.Equal(t, 10, score3, "score = (1+12)-(0+3)") + + // Absent bridgedStats on update -> stored bridged counts left alone. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("update", "cc2", "bafcc2d", + commentRecord("edited3", postURI, postCID, postURI, postCID, nil)))) + _, _, bUp4, bDown4, _, _ := readCommentRow(t, db, uri) + assert.Equal(t, 12, bUp4, "absent bridgedStats leaves stored counts") + assert.Equal(t, 3, bDown4) + + // Equal asOf with DIFFERENT counts -> APPLIED (newer-or-equal guard), native vote + // still preserved. Mirrors the post consumer's equal-asOf case. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("update", "cc2", "bafcc2e", + commentRecord("edited4", postURI, postCID, postURI, postCID, bridgedStatsRecord(20, 2, asOfLate))))) + up5, _, bUp5, bDown5, score5, _ := readCommentRow(t, db, uri) + assert.Equal(t, 1, up5, "native vote preserved") + assert.Equal(t, 20, bUp5, "equal asOf with fresher counts must apply") + assert.Equal(t, 2, bDown5) + assert.Equal(t, 19, score5, "score = (1+20)-(0+2)") +} + +// TestParseRecord_BridgedStats verifies the record parsers tolerate presence/absence +// of bridgedStats without a database. +func TestParseRecord_BridgedStats(t *testing.T) { + withStats, err := parsePostRecord(postRecord("t", "c", bridgedStatsRecord(7, 2, asOfEarly))) + require.NoError(t, err) + require.NotNil(t, withStats.BridgedStats) + assert.Equal(t, 7, withStats.BridgedStats.Upvotes) + assert.Equal(t, 2, withStats.BridgedStats.Downvotes) + assert.Equal(t, asOfEarly, withStats.BridgedStats.AsOf) + + without, err := parsePostRecord(postRecord("t", "c", nil)) + require.NoError(t, err) + assert.Nil(t, without.BridgedStats, "absent bridgedStats parses as nil") + + cWith, err := parseCommentRecord(commentRecord("hi", "at://x/c/1", "cid", "at://x/c/1", "cid", bridgedStatsRecord(3, 1, asOfLate))) + require.NoError(t, err) + require.NotNil(t, cWith.BridgedStats) + assert.Equal(t, 3, cWith.BridgedStats.Upvotes) + + cWithout, err := parseCommentRecord(commentRecord("hi", "at://x/c/1", "cid", "at://x/c/1", "cid", nil)) + require.NoError(t, err) + assert.Nil(t, cWithout.BridgedStats) +} + +// --- edited_at churn (fix 4) --- + +func TestPostConsumer_Update_StatsOnly_EditedAtUnchanged(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/edat1" + + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "edat1", "bafedat1", + postRecord("title", "body", bridgedStatsRecord(5, 1, asOfEarly))))) + + // Stats-only refresh: identical content, newer bridgedStats. A debounced stats + // refresh must NOT mark the post edited. + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("update", "edat1", "bafedat1b", + postRecord("title", "body", bridgedStatsRecord(9, 1, asOfLate))))) + var editedAt *time.Time + require.NoError(t, db.QueryRow(`SELECT edited_at FROM posts WHERE uri=$1`, uri).Scan(&editedAt)) + assert.Nil(t, editedAt, "stats-only update must leave edited_at NULL") + _, _, bUp, _, _, _, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 9, bUp, "bridged counts still refreshed by the stats-only update") + + // A genuine content edit DOES set edited_at. + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("update", "edat1", "bafedat1c", + postRecord("title changed", "body", bridgedStatsRecord(9, 1, asOfLate))))) + require.NoError(t, db.QueryRow(`SELECT edited_at FROM posts WHERE uri=$1`, uri).Scan(&editedAt)) + assert.NotNil(t, editedAt, "content edit must set edited_at") +} + +// --- provenance gate (fix 1a): untrusted repos cannot self-assert bridgedStats --- + +func TestPostConsumer_Create_UntrustedCommunity_BridgedStatsIgnored(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + // Community hosted on a NON-bridge PDS -> provenance gate denies bridgedStats. + insertBridgedCommunityOnPDS(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor, bridgedTestNativePDS) + c := newPostConsumer(t, db) + ctx := context.Background() + uri := "at://" + bridgedTestCommunity + "/social.coves.community.post/unt1" + + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "unt1", "bafunt1", + postRecord("t", "c", bridgedStatsRecord(50, 5, asOfEarly))))) + + // Post indexed normally, but bridgedStats ignored (default-deny). + up, down, bUp, bDown, score, asOf, _, _ := readPostRow(t, db, uri) + assert.Equal(t, 0, up) + assert.Equal(t, 0, down) + assert.Equal(t, 0, bUp, "untrusted community cannot self-assert bridged upvotes") + assert.Equal(t, 0, bDown) + assert.Equal(t, 0, score) + assert.Nil(t, asOf, "no bridged asOf recorded for an ignored aggregate") +} + +func TestCommentConsumer_Create_UntrustedCommenter_BridgedStatsIgnored(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + postURI, postCID := setupCommentThread(t, db) + // Override the commenter to a non-bridge PDS -> provenance gate denies bridgedStats. + insertBridgedUserOnPDS(t, db, bridgedTestCommenter, "brcommenter.test", bridgedTestNativePDS) + cc := newCommentConsumer(db) + ctx := context.Background() + uri := "at://" + bridgedTestCommenter + "/social.coves.community.comment/unt1" + + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("create", "unt1", "bafcunt1", + commentRecord("hi", postURI, postCID, postURI, postCID, bridgedStatsRecord(7, 2, asOfEarly))))) + + up, down, bUp, bDown, score, asOf := readCommentRow(t, db, uri) + assert.Equal(t, 0, up) + assert.Equal(t, 0, down) + assert.Equal(t, 0, bUp, "untrusted commenter cannot self-assert bridged upvotes") + assert.Equal(t, 0, bDown) + assert.Equal(t, 0, score) + assert.Nil(t, asOf) +} + +// --- input hygiene (fix 1b): negative / over-cap aggregates are ignored whole --- + +func TestPostConsumer_Create_BridgedStatsHygiene(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + insertBridgedUser(t, db, bridgedTestAuthor, "brauthor.test") + insertBridgedCommunity(t, db, bridgedTestCommunity, "brcommunity.test", bridgedTestAuthor) + c := newPostConsumer(t, db) + ctx := context.Background() + + // Negative count -> whole aggregate ignored. + negURI := "at://" + bridgedTestCommunity + "/social.coves.community.post/neg1" + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "neg1", "bafneg1", + postRecord("t", "c", bridgedStatsRecord(-5, 0, asOfEarly))))) + _, _, bUp, bDown, score, asOf, _, _ := readPostRow(t, db, negURI) + assert.Equal(t, 0, bUp, "negative bridged count ignored") + assert.Equal(t, 0, bDown) + assert.Equal(t, 0, score) + assert.Nil(t, asOf) + + // Count above maxBridgedCount -> whole aggregate ignored. + capURI := "at://" + bridgedTestCommunity + "/social.coves.community.post/cap1" + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "cap1", "bafcap1", + postRecord("t", "c", bridgedStatsRecord(maxBridgedCount+1, 0, asOfEarly))))) + _, _, bUp2, _, _, asOf2, _, _ := readPostRow(t, db, capURI) + assert.Equal(t, 0, bUp2, "over-cap bridged count ignored") + assert.Nil(t, asOf2) + + // Exactly at the cap -> accepted. + okURI := "at://" + bridgedTestCommunity + "/social.coves.community.post/cap2" + require.NoError(t, c.HandleEvent(ctx, postCommitEvent("create", "cap2", "bafcap2", + postRecord("t", "c", bridgedStatsRecord(maxBridgedCount, 0, asOfEarly))))) + _, _, bUp3, _, _, asOf3, _, _ := readPostRow(t, db, okURI) + assert.Equal(t, maxBridgedCount, bUp3, "count exactly at the cap is accepted") + require.NotNil(t, asOf3) +} + +// --- comment update skips soft-deleted rows (fix 6) --- + +func TestCommentConsumer_Update_SoftDeletedSkipped(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + postURI, postCID := setupCommentThread(t, db) + cc := newCommentConsumer(db) + ctx := context.Background() + uri := "at://" + bridgedTestCommenter + "/social.coves.community.comment/cdel1" + + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("create", "cdel1", "bafcdel1", + commentRecord("original", postURI, postCID, postURI, postCID, bridgedStatsRecord(5, 1, asOfEarly))))) + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("delete", "cdel1", "", nil))) + + // Update on a soft-deleted comment must be skipped (no error, no resurrection, no + // spurious ErrCommentNotFound), leaving the row deleted with content still blanked. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("update", "cdel1", "bafcdel1b", + commentRecord("edited", postURI, postCID, postURI, postCID, bridgedStatsRecord(99, 9, asOfLate))))) + + var deletedAt *time.Time + var content string + var bUp int + require.NoError(t, db.QueryRow(`SELECT deleted_at, content, bridged_upvote_count FROM comments WHERE uri=$1`, uri). + Scan(&deletedAt, &content, &bUp)) + assert.NotNil(t, deletedAt, "comment stays soft-deleted") + assert.Equal(t, "", content, "content stays blanked") + assert.Equal(t, 5, bUp, "bridged counts untouched on soft-deleted row") +} + +// --- comment resurrection score invariant (fix 5) --- + +func TestCommentConsumer_Resurrection_ScoreIncludesSurvivingNativeVotes(t *testing.T) { + db := setupBridgedTestDB(t) + defer func() { _ = db.Close() }() + defer cleanupBridgedTestData(t, db) + cleanupBridgedTestData(t, db) + + postURI, postCID := setupCommentThread(t, db) + cc := newCommentConsumer(db) + vc := newVoteConsumer(db) + ctx := context.Background() + uri := "at://" + bridgedTestCommenter + "/social.coves.community.comment/res1" + + // Create (no bridgedStats), then two native upvotes from distinct voters. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("create", "res1", "bafres1", + commentRecord("hi", postURI, postCID, postURI, postCID, nil)))) + upvote := func(voter, rkey string) { + require.NoError(t, vc.HandleEvent(ctx, &JetstreamEvent{ + Kind: "commit", Did: voter, + Commit: &CommitEvent{Operation: "create", Collection: "social.coves.feed.vote", RKey: rkey, CID: "bafresv", + Record: map[string]interface{}{ + "subject": map[string]interface{}{"uri": uri, "cid": "bafres1"}, + "direction": "up", "createdAt": "2026-02-01T00:00:00Z", + }}, + })) + } + upvote(bridgedTestVoter, "rv1") + upvote(bridgedTestVoter+"2", "rv2") + + up, _, _, _, score, _ := readCommentRow(t, db, uri) + require.Equal(t, 2, up, "two native upvotes recorded") + require.Equal(t, 2, score) + + // Soft-delete, then resurrect (same rkey) WITH bridgedStats. The surviving native + // upvotes must be reflected in the recomputed score. + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("delete", "res1", "", nil))) + require.NoError(t, cc.HandleEvent(ctx, commentCommitEvent("create", "res1", "bafres1b", + commentRecord("back", postURI, postCID, postURI, postCID, bridgedStatsRecord(4, 1, asOfLate))))) + + up2, down2, bUp2, bDown2, score2, _ := readCommentRow(t, db, uri) + assert.Equal(t, 2, up2, "native upvotes survive resurrection") + assert.Equal(t, 0, down2) + assert.Equal(t, 4, bUp2) + assert.Equal(t, 1, bDown2) + assert.Equal(t, 5, score2, "score = (2+4)-(0+1) includes surviving native votes") +} diff --git a/internal/atproto/jetstream/comment_consumer.go b/internal/atproto/jetstream/comment_consumer.go index 72f4819..9cee809 100644 --- a/internal/atproto/jetstream/comment_consumer.go +++ b/internal/atproto/jetstream/comment_consumer.go @@ -32,17 +32,59 @@ const ( type CommentEventConsumer struct { commentRepo comments.Repository db *sql.DB // Direct DB access for atomic count updates + // bridgeTrust gates whether a comment's user repo may assert bridgedStats. + // nil means default-deny (bridgedStats are ignored for every comment). + bridgeTrust *BridgeTrust +} + +// CommentEventConsumerOption configures optional CommentEventConsumer behaviour. +type CommentEventConsumerOption func(*CommentEventConsumer) + +// WithCommentBridgeTrust installs the provenance gate that decides which user repos may +// assert bridgedStats on their comments. Without it, bridgedStats are default-denied. +func WithCommentBridgeTrust(bt *BridgeTrust) CommentEventConsumerOption { + return func(c *CommentEventConsumer) { c.bridgeTrust = bt } } // NewCommentEventConsumer creates a new Jetstream consumer for comment events func NewCommentEventConsumer( commentRepo comments.Repository, db *sql.DB, + opts ...CommentEventConsumerOption, ) *CommentEventConsumer { - return &CommentEventConsumer{ + c := &CommentEventConsumer{ commentRepo: commentRepo, db: db, } + for _, opt := range opts { + opt(c) + } + return c +} + +// bridgeStatsAllowedForRepo reports whether the given comment repo (a user DID) is a +// trusted bridge, i.e. whether its records may assert bridgedStats. It resolves the +// repo's PDS host from the already-indexed users row (users.pds_url, populated from +// identity resolution at user creation) and checks it against the trust allowlist. +// Default-deny: any lookup failure — including the user not being indexed yet — means +// "not trusted", so bridgedStats are ignored. The accepted consequence: a bridged +// comment that arrives before its (bridge) author is indexed has its bridgedStats +// dropped at create; they are folded in on the next record edit once the author exists. +func (c *CommentEventConsumer) bridgeStatsAllowedForRepo(ctx context.Context, repoDID string) bool { + if c.bridgeTrust == nil { + return false + } + var pdsURL string + err := c.db.QueryRowContext(ctx, `SELECT pds_url FROM users WHERE did = $1`, repoDID).Scan(&pdsURL) + if err != nil { + if err == sql.ErrNoRows { + log.Printf("debug: ignoring bridgedStats from repo %s (user not indexed; provenance unverifiable)", repoDID) + } else { + log.Printf("Warning: bridgedStats provenance check failed for repo %s: %v", repoDID, err) + } + return false + } + return c.bridgeTrust.TrustsPDS(pdsURL) } // HandleEvent processes a Jetstream event for comment records @@ -124,6 +166,26 @@ func (c *CommentEventConsumer) createComment(ctx context.Context, repoDID string IndexedAt: time.Now(), } + // Apply bridge-asserted origin-platform vote aggregates if the record carries them, + // the user repo is a trusted bridge (provenance gate, default-deny), and the + // aggregate passes input hygiene. At create there are no native votes, so the + // inclusive score is the bridged delta. bridgedStats is optional (absent for + // natively-authored comments). Counts + asOf are applied atomically: an unparseable + // or out-of-hygiene aggregate is ignored WHOLE, never leaving counts with a NULL + // asOf (which would defeat the update-path regression guard). + if commentRecord.BridgedStats != nil { + if c.bridgeStatsAllowedForRepo(ctx, repoDID) { + if up, down, asOf, ok := validatedBridgedStats(commentRecord.BridgedStats, uri); ok { + comment.BridgedUpvoteCount = up + comment.BridgedDownvoteCount = down + comment.BridgedStatsAsOf = &asOf + comment.Score = up - down + } + } else { + log.Printf("debug: ignoring bridgedStats on comment %s from untrusted repo %s (not a trusted bridge PDS)", uri, repoDID) + } + } + // Atomically: Index comment + Update parent counts if err := c.indexCommentAndUpdateCounts(ctx, comment); err != nil { return fmt.Errorf("failed to index comment and update counts: %w", err) @@ -133,7 +195,13 @@ func (c *CommentEventConsumer) createComment(ctx context.Context, repoDID string return nil } -// updateComment updates an existing comment's content fields +// updateComment updates an existing comment's content fields. +// +// Like updatePost, this is idempotent and error-return means log-and-drop (the +// connector tracks no cursor and live-tails Jetstream, so a returned error is NOT +// replayed): the folded bridged counts only self-heal on the bridge's next record +// edit. We therefore skip benign no-ops (missing row, soft-deleted row) cleanly and +// reserve errors for transient infra faults. func (c *CommentEventConsumer) updateComment(ctx context.Context, repoDID string, commit *CommitEvent) error { if commit.Record == nil { return fmt.Errorf("comment update event missing record data") @@ -154,27 +222,49 @@ func (c *CommentEventConsumer) updateComment(ctx context.Context, repoDID string // Build AT-URI for the comment being updated uri := fmt.Sprintf("at://%s/social.coves.community.comment/%s", repoDID, commit.RKey) - // Fetch existing comment to validate threading references are immutable - existingComment, err := c.commentRepo.GetByURI(ctx, uri) + // Load the raw stored columns needed to enforce threading immutability, skip + // soft-deleted rows, and run the bridgedStats regression guard. This is a dedicated + // consumer-only query (mirroring post_consumer's inline SELECT): it reads the RAW + // native and bridged columns, deliberately NOT the folded display counts that + // GetByURI now returns, so the guard reasons about the true separate aggregates. + var ( + storedRootURI, storedRootCID string + storedParentURI, storedParentCID string + storedDeletedAt *time.Time + storedAsOf *time.Time + ) + err = c.db.QueryRowContext(ctx, + `SELECT root_uri, root_cid, parent_uri, parent_cid, deleted_at, bridged_stats_as_of + FROM comments WHERE uri = $1`, uri, + ).Scan(&storedRootURI, &storedRootCID, &storedParentURI, &storedParentCID, &storedDeletedAt, &storedAsOf) + if err == sql.ErrNoRows { + // Comment not indexed yet: its CREATE event will index it when it arrives. + log.Printf("Update event for non-indexed comment: %s (will be indexed on CREATE)", uri) + return nil + } if err != nil { - if err == comments.ErrCommentNotFound { - // Comment doesn't exist yet - might arrive out of order - log.Printf("Warning: Update event for non-existent comment: %s (will be indexed on CREATE)", uri) - return nil - } - return fmt.Errorf("failed to get existing comment for validation: %w", err) + return fmt.Errorf("failed to load stored comment for update: %w", err) + } + + // Skip soft-deleted rows. A moderator-removed (or author-deleted) comment must not be + // resurrected by a debounced stats edit; the repo Update's deleted_at IS NULL guard + // would otherwise match no rows and surface as a spurious ErrCommentNotFound failure + // recurring on every stats refresh (mirrors post_consumer's soft-deleted skip). + if storedDeletedAt != nil { + log.Printf("Update event for soft-deleted comment: %s (skipping)", uri) + return nil } // SECURITY: Threading references are IMMUTABLE after creation // Reject updates that attempt to change root/parent (prevents thread hijacking) - if existingComment.RootURI != commentRecord.Reply.Root.URI || - existingComment.RootCID != commentRecord.Reply.Root.CID || - existingComment.ParentURI != commentRecord.Reply.Parent.URI || - existingComment.ParentCID != commentRecord.Reply.Parent.CID { + if storedRootURI != commentRecord.Reply.Root.URI || + storedRootCID != commentRecord.Reply.Root.CID || + storedParentURI != commentRecord.Reply.Parent.URI || + storedParentCID != commentRecord.Reply.Parent.CID { log.Printf("🚨 SECURITY: Rejecting comment update - threading references are immutable: %s", uri) - log.Printf(" Existing root: %s (CID: %s)", existingComment.RootURI, existingComment.RootCID) + log.Printf(" Existing root: %s (CID: %s)", storedRootURI, storedRootCID) log.Printf(" Incoming root: %s (CID: %s)", commentRecord.Reply.Root.URI, commentRecord.Reply.Root.CID) - log.Printf(" Existing parent: %s (CID: %s)", existingComment.ParentURI, existingComment.ParentCID) + log.Printf(" Existing parent: %s (CID: %s)", storedParentURI, storedParentCID) log.Printf(" Incoming parent: %s (CID: %s)", commentRecord.Reply.Parent.URI, commentRecord.Reply.Parent.CID) return fmt.Errorf("comment threading references cannot be changed after creation") } @@ -185,15 +275,46 @@ func (c *CommentEventConsumer) updateComment(ctx context.Context, repoDID string return fmt.Errorf("failed to serialize optional fields: %w", err) } - // Build comment update entity (preserves vote counts and created_at) + // Decide the candidate bridged aggregate handed to repo.Update. It is applied only + // when the record carried bridgedStats AND the user repo is a trusted bridge + // (provenance gate, default-deny) AND the aggregate passes input hygiene AND its asOf + // parses. Otherwise we pass a nil incoming asOf, and the atomic SQL guard in + // repo.Update leaves the stored bridged columns untouched. The newer-or-equal + // regression comparison happens ATOMICALLY inside the UPDATE; storedAsOf is read here + // only to log the strictly-older case. + var ( + incomingUp, incomingDn int + incomingAsOf *time.Time + ) + if commentRecord.BridgedStats != nil { + if c.bridgeStatsAllowedForRepo(ctx, repoDID) { + if up, down, asOf, ok := validatedBridgedStats(commentRecord.BridgedStats, uri); ok { + incomingUp, incomingDn, incomingAsOf = up, down, &asOf + if storedAsOf != nil && asOf.Before(*storedAsOf) { + log.Printf("debug: ignoring strictly-older bridgedStats for %s (incoming asOf %s < stored %s)", + uri, asOf.Format(time.RFC3339), storedAsOf.Format(time.RFC3339)) + } + } + } else { + log.Printf("debug: ignoring bridgedStats on comment %s from untrusted repo %s (not a trusted bridge PDS)", uri, repoDID) + } + } + + // Build comment update entity (preserves native vote counts and created_at). The + // Bridged* fields carry the INCOMING candidate; repo.Update applies them atomically + // only when incomingAsOf is non-nil and newer-or-equal to the stored asOf, and always + // recomputes the inclusive score from live native counts so concurrent votes survive. comment := &comments.Comment{ - URI: uri, - CID: commit.CID, - Content: commentRecord.Content, - ContentFacets: facetsJSON, - Embed: embedJSON, - ContentLabels: labelsJSON, - Langs: commentRecord.Langs, + URI: uri, + CID: commit.CID, + Content: commentRecord.Content, + ContentFacets: facetsJSON, + Embed: embedJSON, + ContentLabels: labelsJSON, + Langs: commentRecord.Langs, + BridgedUpvoteCount: incomingUp, + BridgedDownvoteCount: incomingDn, + BridgedStatsAsOf: incomingAsOf, } // Update the comment in repository @@ -201,7 +322,11 @@ func (c *CommentEventConsumer) updateComment(ctx context.Context, repoDID string return fmt.Errorf("failed to update comment: %w", err) } - log.Printf("✓ Updated comment: %s", uri) + if incomingAsOf != nil { + log.Printf("✓ Updated comment: %s (bridgedStats candidate applied if newer-or-equal: up=%d down=%d)", uri, incomingUp, incomingDn) + } else { + log.Printf("✓ Updated comment: %s", uri) + } return nil } @@ -296,8 +421,17 @@ func (c *CommentEventConsumer) indexCommentAndUpdateCounts(ctx context.Context, deleted_at = NULL, deletion_reason = NULL, deleted_by = NULL, - reply_count = 0 - WHERE id = $14 + reply_count = 0, + bridged_upvote_count = $14, + bridged_downvote_count = $15, + bridged_stats_as_of = $16, + -- Recompute the inclusive score in SQL from the SURVIVING native counts + -- (upvote_count/downvote_count are intentionally not reset on resurrect) + -- plus the incoming bridged values, upholding migration 031's invariant + -- score = (up+bUp) - (down+bDown). Using comment.Score here would have + -- written only the bridged delta and dropped the retained native votes. + score = upvote_count + $14 - downvote_count - $15 + WHERE id = $17 ` _, err = tx.ExecContext( @@ -315,6 +449,9 @@ func (c *CommentEventConsumer) indexCommentAndUpdateCounts(ctx context.Context, pq.Array(comment.Langs), comment.CreatedAt, time.Now(), + comment.BridgedUpvoteCount, + comment.BridgedDownvoteCount, + comment.BridgedStatsAsOf, commentID, ) if err != nil { @@ -330,12 +467,14 @@ func (c *CommentEventConsumer) indexCommentAndUpdateCounts(ctx context.Context, uri, cid, rkey, commenter_did, root_uri, root_cid, parent_uri, parent_cid, content, content_facets, embed, content_labels, langs, - created_at, indexed_at + created_at, indexed_at, + bridged_upvote_count, bridged_downvote_count, bridged_stats_as_of, score ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, - $14, $15 + $14, $15, + $16, $17, $18, $19 ) ON CONFLICT (uri) DO NOTHING RETURNING id @@ -347,6 +486,7 @@ func (c *CommentEventConsumer) indexCommentAndUpdateCounts(ctx context.Context, comment.RootURI, comment.RootCID, comment.ParentURI, comment.ParentCID, comment.Content, comment.ContentFacets, comment.Embed, comment.ContentLabels, pq.Array(comment.Langs), comment.CreatedAt, time.Now(), + comment.BridgedUpvoteCount, comment.BridgedDownvoteCount, comment.BridgedStatsAsOf, comment.Score, ).Scan(&commentID) if err == sql.ErrNoRows { // ON CONFLICT triggered - comment was inserted by concurrent process @@ -634,14 +774,15 @@ func validateATURI(uri string) error { // CommentRecordFromJetstream represents a comment record as received from Jetstream // Matches social.coves.community.comment lexicon type CommentRecordFromJetstream struct { - Labels interface{} `json:"labels,omitempty"` - Embed map[string]interface{} `json:"embed,omitempty"` - Reply ReplyRefFromJetstream `json:"reply"` - Type string `json:"$type"` - Content string `json:"content"` - CreatedAt string `json:"createdAt"` - Facets []interface{} `json:"facets,omitempty"` - Langs []string `json:"langs,omitempty"` + Labels interface{} `json:"labels,omitempty"` + Embed map[string]interface{} `json:"embed,omitempty"` + BridgedStats *BridgedStatsFromJetstream `json:"bridgedStats,omitempty"` + Reply ReplyRefFromJetstream `json:"reply"` + Type string `json:"$type"` + Content string `json:"content"` + CreatedAt string `json:"createdAt"` + Facets []interface{} `json:"facets,omitempty"` + Langs []string `json:"langs,omitempty"` } // ReplyRefFromJetstream represents the threading structure diff --git a/internal/atproto/jetstream/post_consumer.go b/internal/atproto/jetstream/post_consumer.go index 7632205..9e8d073 100644 --- a/internal/atproto/jetstream/post_consumer.go +++ b/internal/atproto/jetstream/post_consumer.go @@ -14,13 +14,24 @@ import ( ) // PostEventConsumer consumes post-related events from Jetstream -// Handles CREATE and DELETE operations for social.coves.community.post -// UPDATE handler will be added when that feature is implemented +// Handles CREATE, UPDATE, and DELETE operations for social.coves.community.post type PostEventConsumer struct { postRepo posts.Repository communityRepo communities.Repository userService users.UserService db *sql.DB // Direct DB access for atomic count reconciliation + // bridgeTrust gates whether a post's community repo may assert bridgedStats. + // nil means default-deny (bridgedStats are ignored for every post). + bridgeTrust *BridgeTrust +} + +// PostEventConsumerOption configures optional PostEventConsumer behaviour. +type PostEventConsumerOption func(*PostEventConsumer) + +// WithPostBridgeTrust installs the provenance gate that decides which community repos +// may assert bridgedStats on their posts. Without it, bridgedStats are default-denied. +func WithPostBridgeTrust(bt *BridgeTrust) PostEventConsumerOption { + return func(c *PostEventConsumer) { c.bridgeTrust = bt } } // NewPostEventConsumer creates a new Jetstream consumer for post events @@ -29,17 +40,22 @@ func NewPostEventConsumer( communityRepo communities.Repository, userService users.UserService, db *sql.DB, + opts ...PostEventConsumerOption, ) *PostEventConsumer { - return &PostEventConsumer{ + c := &PostEventConsumer{ postRepo: postRepo, communityRepo: communityRepo, userService: userService, db: db, } + for _, opt := range opts { + opt(c) + } + return c } // HandleEvent processes a Jetstream event for post records -// Handles CREATE and DELETE operations - UPDATE deferred until that feature exists +// Handles CREATE, UPDATE, and DELETE operations func (c *PostEventConsumer) HandleEvent(ctx context.Context, event *JetstreamEvent) error { // We only care about commit events for post records if event.Kind != "commit" || event.Commit == nil { @@ -53,12 +69,14 @@ func (c *PostEventConsumer) HandleEvent(ctx context.Context, event *JetstreamEve switch commit.Operation { case "create": return c.createPost(ctx, event.Did, commit) + case "update": + return c.updatePost(ctx, event.Did, commit) case "delete": return c.deletePost(ctx, event.Did, commit) } } - // Silently ignore other operations (update) and other collections + // Silently ignore other operations and other collections return nil } @@ -74,9 +92,11 @@ func (c *PostEventConsumer) createPost(ctx context.Context, repoDID string, comm return fmt.Errorf("failed to parse post record: %w", err) } - // SECURITY: Validate this is a legitimate post event - if err := c.validatePostEvent(ctx, repoDID, postRecord); err != nil { - log.Printf("🚨 SECURITY: Rejecting post event: %v", err) + // SECURITY: Validate this is a legitimate post event. Returns the community row so + // we can check bridgedStats provenance against its resolved PDS host. + community, err := c.validatePostEvent(ctx, repoDID, postRecord) + if err != nil { + logPostValidationRejection("post create", err) return err } @@ -103,13 +123,33 @@ func (c *PostEventConsumer) createPost(ctx context.Context, repoDID string, comm Content: postRecord.Content, CreatedAt: createdAt, IndexedAt: time.Now(), - // Stats remain at 0 (no votes yet) + // Native stats remain at 0 (no native votes yet); bridged stats applied below. UpvoteCount: 0, DownvoteCount: 0, Score: 0, CommentCount: 0, } + // Apply bridge-asserted origin-platform vote aggregates if the record carries them, + // the community repo is a trusted bridge (provenance gate, default-deny), and the + // aggregate passes input hygiene. At create there are no native votes, so the + // inclusive score is simply the bridged delta. bridgedStats is optional (absent for + // natively-authored posts). The counts + asOf are applied atomically: an unparseable + // or out-of-hygiene aggregate is ignored WHOLE, never leaving counts with a NULL + // asOf (which would defeat the update-path regression guard). + if postRecord.BridgedStats != nil { + if c.bridgeTrust.TrustsPDS(community.PDSURL) { + if up, down, asOf, ok := validatedBridgedStats(postRecord.BridgedStats, uri); ok { + post.BridgedUpvoteCount = up + post.BridgedDownvoteCount = down + post.BridgedStatsAsOf = &asOf + post.Score = up - down + } + } else { + log.Printf("debug: ignoring bridgedStats on post %s from untrusted repo %s (not a trusted bridge PDS)", uri, repoDID) + } + } + // Serialize JSON fields (facets, embed, labels) // Return error if any non-empty field fails to serialize (prevents silent data loss) if postRecord.Facets != nil { @@ -165,6 +205,210 @@ func (c *PostEventConsumer) deletePost(ctx context.Context, repoDID string, comm return nil } +// updatePost handles post record update events from Jetstream. +// +// Posts previously ignored updates; the bridge now edits post records (content and +// especially the refreshed bridgedStats aggregate) via debounced record updates, so +// we must fold those into the index. Every branch below is idempotent, which matters +// because the connector logs-and-drops on error WITHOUT tracking a cursor (it live- +// tails Jetstream): a returned error is NOT retried or replayed, so we only return an +// error for genuinely transient infra faults and otherwise skip benign no-ops cleanly. +// The accepted consequence for stats: if a bridgedStats update errors out, the folded +// counts stay stale until the bridge next edits the record (which it does on every +// stats refresh), so the desync is self-healing rather than permanent. +// +// Security mirrors createPost (repoDID must equal record.community; community and +// author must exist). We additionally reject reassignment: an update may not move a +// post to a different community or author. Reassignment, a missing stored row, and a +// soft-deleted stored row are all skipped (logged, no error). +func (c *PostEventConsumer) updatePost(ctx context.Context, repoDID string, commit *CommitEvent) error { + if commit.Record == nil { + return fmt.Errorf("post update event missing record data") + } + + postRecord, err := parsePostRecord(commit.Record) + if err != nil { + return fmt.Errorf("failed to parse post record: %w", err) + } + + // SECURITY: identical validation to create (repo == community, community/author exist). + community, err := c.validatePostEvent(ctx, repoDID, postRecord) + if err != nil { + logPostValidationRejection("post update", err) + return err + } + + uri := fmt.Sprintf("at://%s/social.coves.community.post/%s", repoDID, commit.RKey) + + // Fetch the stored row so we can enforce immutability and run the asOf regression guard. + var ( + storedID int64 + storedCommunityDID string + storedAuthorDID string + storedDeletedAt *time.Time + storedAsOf *time.Time + ) + err = c.db.QueryRowContext(ctx, + `SELECT id, community_did, author_did, deleted_at, bridged_stats_as_of FROM posts WHERE uri = $1`, + uri, + ).Scan(&storedID, &storedCommunityDID, &storedAuthorDID, &storedDeletedAt, &storedAsOf) + if err == sql.ErrNoRows { + // Not indexed yet (out-of-order delivery). Jetstream will replay CREATE; skip. + log.Printf("Update event for non-indexed post: %s (will be indexed on CREATE)", uri) + return nil + } + if err != nil { + return fmt.Errorf("failed to load stored post for update: %w", err) + } + + // Skip soft-deleted rows: a deleted post should not be resurrected by an edit. + if storedDeletedAt != nil { + log.Printf("Update event for soft-deleted post: %s (skipping)", uri) + return nil + } + + // SECURITY: community and author are immutable. Reassignment is rejected (skipped). + if storedCommunityDID != postRecord.Community || storedAuthorDID != postRecord.Author { + log.Printf("🚨 SECURITY: Rejecting post update - community/author reassignment is not allowed: %s (stored community=%s author=%s; incoming community=%s author=%s)", + uri, storedCommunityDID, storedAuthorDID, postRecord.Community, postRecord.Author) + return nil + } + + // Serialize optional JSON content fields (return on failure to avoid silent data loss). + var facetsJSON, embedJSON, labelsJSON sql.NullString + if postRecord.Facets != nil { + b, marshalErr := json.Marshal(postRecord.Facets) + if marshalErr != nil { + return fmt.Errorf("failed to serialize facets: %w", marshalErr) + } + facetsJSON.String, facetsJSON.Valid = string(b), true + } + if postRecord.Embed != nil { + b, marshalErr := json.Marshal(postRecord.Embed) + if marshalErr != nil { + return fmt.Errorf("failed to serialize embed: %w", marshalErr) + } + embedJSON.String, embedJSON.Valid = string(b), true + } + if postRecord.Labels != nil { + b, marshalErr := json.Marshal(postRecord.Labels) + if marshalErr != nil { + return fmt.Errorf("failed to serialize labels: %w", marshalErr) + } + labelsJSON.String, labelsJSON.Valid = string(b), true + } + + // Decide the candidate bridged aggregate to hand to the atomic UPDATE. It is applied + // only when the record carried bridgedStats AND the community repo is a trusted + // bridge (provenance gate, default-deny) AND the aggregate passes input hygiene + // (non-negative, within the magnitude cap) AND its asOf parses. Otherwise we pass a + // NULL incoming asOf, which makes the SQL guard below leave the stored bridged + // columns untouched. The actual newer-or-equal regression comparison is done + // ATOMICALLY inside the UPDATE (see the CASE expressions) rather than read-here / + // write-later, so it cannot race a concurrent write; storedAsOf is read only to log + // the strictly-older case. + var ( + incomingUp, incomingDown int + incomingAsOf *time.Time + ) + if postRecord.BridgedStats != nil { + if c.bridgeTrust.TrustsPDS(community.PDSURL) { + if up, down, asOf, ok := validatedBridgedStats(postRecord.BridgedStats, uri); ok { + incomingUp, incomingDown, incomingAsOf = up, down, &asOf + // Best-effort log only (the write is authoritative and atomic): a + // strictly-older asOf is dropped by the SQL guard. Kept at debug because + // the bridge re-sends the same asOf on every content edit, so this is + // noise, not an anomaly. + if storedAsOf != nil && asOf.Before(*storedAsOf) { + log.Printf("debug: ignoring strictly-older bridgedStats for %s (incoming asOf %s < stored %s)", + uri, asOf.Format(time.RFC3339), storedAsOf.Format(time.RFC3339)) + } + } + } else { + log.Printf("debug: ignoring bridgedStats on post %s from untrusted repo %s (not a trusted bridge PDS)", uri, repoDID) + } + } + + // Single atomic UPDATE. edited_at is bumped only when content actually changed (so a + // debounced stats-only refresh does not mark the post edited). The bridged columns + // and the inclusive score move together via a shared applies-guard: apply the + // incoming counts only when an incoming asOf is present and is newer-or-equal to the + // stored one (NULL stored => first application). score is always recomputed from the + // LIVE native counts plus whichever bridged counts win, so concurrent native votes + // are never clobbered. $10 is the incoming asOf (NULL => no bridged change); the + // stored asOf is read directly from the row, keeping the compare atomic. + updateQuery := ` + UPDATE posts + SET + cid = $2, + title = $3, + content = $4, + content_facets = $5, + embed = $6, + content_labels = $7, + edited_at = CASE + WHEN title IS DISTINCT FROM $3 OR content IS DISTINCT FROM $4 + OR content_facets IS DISTINCT FROM $5 OR embed IS DISTINCT FROM $6 + OR content_labels IS DISTINCT FROM $7 + THEN NOW() ELSE edited_at END, + bridged_upvote_count = CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $8 ELSE bridged_upvote_count END, + bridged_downvote_count = CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $9 ELSE bridged_downvote_count END, + bridged_stats_as_of = CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $10 ELSE bridged_stats_as_of END, + score = upvote_count - downvote_count + + CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $8 ELSE bridged_upvote_count END + - CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $9 ELSE bridged_downvote_count END + WHERE id = $1 AND deleted_at IS NULL + ` + result, err := c.db.ExecContext(ctx, updateQuery, + storedID, commit.CID, postRecord.Title, postRecord.Content, + facetsJSON, embedJSON, labelsJSON, + incomingUp, incomingDown, incomingAsOf, + ) + if err != nil { + return fmt.Errorf("failed to update post: %w", err) + } + + // A post can be soft-deleted between the load above and this UPDATE; the + // deleted_at IS NULL guard then matches no rows. Report that as a skip instead of + // falsely logging a successful update (mirrors vote_consumer's RowsAffected check). + rowsAffected, err := result.RowsAffected() + if err != nil { + return fmt.Errorf("failed to check post update result: %w", err) + } + if rowsAffected == 0 { + log.Printf("Update event for post that vanished/was deleted between load and update: %s (skipping)", uri) + return nil + } + + if incomingAsOf != nil { + log.Printf("✓ Updated post: %s (bridgedStats candidate applied if newer-or-equal: up=%d down=%d)", uri, incomingUp, incomingDown) + } else { + log.Printf("✓ Updated post: %s", uri) + } + return nil +} + +// parseBridgedAsOf parses a bridgedStats.asOf timestamp, logging (and returning the +// error) on failure so callers can decide to skip applying the aggregate. +func parseBridgedAsOf(asOf, uri string) (time.Time, error) { + t, err := time.Parse(time.RFC3339, asOf) + if err != nil { + log.Printf("Warning: failed to parse bridgedStats.asOf %q for %s: %v", asOf, uri, err) + return time.Time{}, err + } + return t, nil +} + // indexPostAndReconcileCounts atomically indexes a post and reconciles comment counts // This fixes the race condition where comments arrive before their parent post func (c *PostEventConsumer) indexPostAndReconcileCounts(ctx context.Context, post *posts.Post) error { @@ -200,11 +444,13 @@ func (c *PostEventConsumer) indexPostAndReconcileCounts(ctx context.Context, pos INSERT INTO posts ( uri, cid, rkey, author_did, community_did, title, content, content_facets, embed, content_labels, - created_at, indexed_at + created_at, indexed_at, + bridged_upvote_count, bridged_downvote_count, bridged_stats_as_of, score ) VALUES ( $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, - $11, NOW() + $11, NOW(), + $12, $13, $14, $15 ) ON CONFLICT (uri) DO NOTHING RETURNING id @@ -216,6 +462,7 @@ func (c *PostEventConsumer) indexPostAndReconcileCounts(ctx context.Context, pos post.URI, post.CID, post.RKey, post.AuthorDID, post.CommunityDID, post.Title, post.Content, facetsJSON, embedJSON, labelsJSON, post.CreatedAt, + post.BridgedUpvoteCount, post.BridgedDownvoteCount, post.BridgedStatsAsOf, post.Score, ).Scan(&postID) // If no rows returned, post already exists (idempotent - OK for Jetstream replays) @@ -268,9 +515,28 @@ func (c *PostEventConsumer) indexPostAndReconcileCounts(ctx context.Context, pos return nil } -// validatePostEvent performs security validation on post events -// This prevents malicious actors from indexing fake posts -func (c *PostEventConsumer) validatePostEvent(ctx context.Context, repoDID string, post *PostRecordFromJetstream) error { +// errValidationInfra marks a post-validation failure caused by an infrastructure fault +// (e.g. a DB error while checking that the community or author exists) rather than a +// policy rejection. The two are logged differently: policy rejections are security +// events (🚨), infra faults are plain operational errors that must NOT masquerade as +// an attack in the logs. +var errValidationInfra = errors.New("validation infrastructure error") + +// logPostValidationRejection logs a validatePostEvent failure, distinguishing genuine +// policy rejections (security-relevant) from infrastructure faults (operational). See +// errValidationInfra. +func logPostValidationRejection(op string, err error) { + if errors.Is(err, errValidationInfra) { + log.Printf("Error: %s could not be validated (infrastructure fault, not a rejection): %v", op, err) + return + } + log.Printf("🚨 SECURITY: Rejecting %s: %v", op, err) +} + +// validatePostEvent performs security validation on post events and, on success, +// returns the community row (whose resolved PDS host drives the bridgedStats +// provenance gate). This prevents malicious actors from indexing fake posts. +func (c *PostEventConsumer) validatePostEvent(ctx context.Context, repoDID string, post *PostRecordFromJetstream) (*communities.Community, error) { // CRITICAL SECURITY CHECK: // Posts MUST come from community repositories, not user repositories // This prevents users from creating posts that appear to be from communities they don't control @@ -284,59 +550,70 @@ func (c *PostEventConsumer) validatePostEvent(ctx context.Context, repoDID strin // - We verify event.Did (repo owner) == post.community (claimed community) // - Reject if mismatch if repoDID != post.Community { - return fmt.Errorf("repository DID (%s) doesn't match community DID (%s) - posts must come from community repos", + return nil, fmt.Errorf("repository DID (%s) doesn't match community DID (%s) - posts must come from community repos", repoDID, post.Community) } - // CRITICAL: Verify community exists in AppView - // Posts MUST reference valid communities (enforced by FK constraint) - // If community isn't indexed yet, we must reject the post - // Jetstream will replay events, so the post will be indexed once community is ready - _, err := c.communityRepo.GetByDID(ctx, post.Community) + // CRITICAL: Verify community exists in AppView. + // Posts MUST reference valid communities (enforced by FK constraint). If the + // community isn't indexed yet we reject; because the connector does not track a + // cursor, the post is only re-indexed if the record is re-emitted (which the bridge + // does on edits) rather than automatically replayed. + community, err := c.communityRepo.GetByDID(ctx, post.Community) if err != nil { if communities.IsNotFound(err) { - // Reject - community must be indexed before posts - // This maintains referential integrity and prevents orphaned posts - return fmt.Errorf("community not found: %s - cannot index post before community", post.Community) + // Policy rejection - community must be indexed before posts. + return nil, fmt.Errorf("community not found: %s - cannot index post before community", post.Community) } - // Database error or other issue - return fmt.Errorf("failed to verify community exists: %w", err) + // Infrastructure fault (DB error): not an attack. Tag it so the caller logs it + // as an operational error, not a 🚨 rejection. + return nil, fmt.Errorf("%w: failed to verify community exists: %v", errValidationInfra, err) } - // CRITICAL: Verify author exists in AppView - // Every post MUST have a valid author (enforced by FK constraint) - // Even though posts live in community repos, they belong to specific authors - // If author isn't indexed yet, we must reject the post + // CRITICAL: Verify author exists in AppView. + // Every post MUST have a valid author (enforced by FK constraint). Even though posts + // live in community repos, they belong to specific authors. _, err = c.userService.GetUserByDID(ctx, post.Author) if err != nil { // Use proper error type checking with errors.Is() if errors.Is(err, users.ErrUserNotFound) { - // Reject - author must be indexed before posts - // This maintains referential integrity and prevents orphaned posts - return fmt.Errorf("author not found: %s - cannot index post before author", post.Author) + // Policy rejection - author must be indexed before posts. + return nil, fmt.Errorf("author not found: %s - cannot index post before author", post.Author) } - // Database error or other issue - return fmt.Errorf("failed to verify author exists: %w", err) + // Infrastructure fault (DB error): not an attack. + return nil, fmt.Errorf("%w: failed to verify author exists: %v", errValidationInfra, err) } - return nil + return community, nil } // PostRecordFromJetstream represents a post record as received from Jetstream // Matches the structure written to PDS via social.coves.community.post type PostRecordFromJetstream struct { - OriginalAuthor interface{} `json:"originalAuthor,omitempty"` - FederatedFrom interface{} `json:"federatedFrom,omitempty"` - Location interface{} `json:"location,omitempty"` - Title *string `json:"title,omitempty"` - Content *string `json:"content,omitempty"` - Embed map[string]interface{} `json:"embed,omitempty"` - Labels *posts.SelfLabels `json:"labels,omitempty"` - Type string `json:"$type"` - Community string `json:"community"` - Author string `json:"author"` - CreatedAt string `json:"createdAt"` - Facets []interface{} `json:"facets,omitempty"` + OriginalAuthor interface{} `json:"originalAuthor,omitempty"` + FederatedFrom interface{} `json:"federatedFrom,omitempty"` + Location interface{} `json:"location,omitempty"` + Title *string `json:"title,omitempty"` + Content *string `json:"content,omitempty"` + Embed map[string]interface{} `json:"embed,omitempty"` + Labels *posts.SelfLabels `json:"labels,omitempty"` + BridgedStats *BridgedStatsFromJetstream `json:"bridgedStats,omitempty"` + Type string `json:"$type"` + Community string `json:"community"` + Author string `json:"author"` + CreatedAt string `json:"createdAt"` + Facets []interface{} `json:"facets,omitempty"` +} + +// BridgedStatsFromJetstream is the bridge-asserted aggregate of origin-platform +// votes carried on federated/bridged post and comment records (social.coves +// community.post / community.comment #bridgedStats). A nil pointer means the +// record carried no bridgedStats, which callers treat as "leave stored counts +// alone" rather than "reset to zero". +type BridgedStatsFromJetstream struct { + Upvotes int `json:"upvotes"` + Downvotes int `json:"downvotes"` + AsOf string `json:"asOf"` } // parsePostRecord converts a raw Jetstream record map to a PostRecordFromJetstream diff --git a/internal/atproto/jetstream/vote_consumer.go b/internal/atproto/jetstream/vote_consumer.go index cf690c9..01054cf 100644 --- a/internal/atproto/jetstream/vote_consumer.go +++ b/internal/atproto/jetstream/vote_consumer.go @@ -185,15 +185,15 @@ func (c *VoteEventConsumer) indexVoteAndUpdateCounts(ctx context.Context, vote * var decrementQuery string if existingDirection.String == "up" { if collection == "social.coves.community.post" { - decrementQuery = `UPDATE posts SET upvote_count = GREATEST(0, upvote_count - 1), score = upvote_count - 1 - downvote_count WHERE uri = $1 AND deleted_at IS NULL` + decrementQuery = `UPDATE posts SET upvote_count = GREATEST(0, upvote_count - 1), score = GREATEST(0, upvote_count - 1) - downvote_count + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL` } else if collection == "social.coves.community.comment" { - decrementQuery = `UPDATE comments SET upvote_count = GREATEST(0, upvote_count - 1), score = upvote_count - 1 - downvote_count WHERE uri = $1 AND deleted_at IS NULL` + decrementQuery = `UPDATE comments SET upvote_count = GREATEST(0, upvote_count - 1), score = GREATEST(0, upvote_count - 1) - downvote_count + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL` } } else { if collection == "social.coves.community.post" { - decrementQuery = `UPDATE posts SET downvote_count = GREATEST(0, downvote_count - 1), score = upvote_count - (downvote_count - 1) WHERE uri = $1 AND deleted_at IS NULL` + decrementQuery = `UPDATE posts SET downvote_count = GREATEST(0, downvote_count - 1), score = upvote_count - GREATEST(0, downvote_count - 1) + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL` } else if collection == "social.coves.community.comment" { - decrementQuery = `UPDATE comments SET downvote_count = GREATEST(0, downvote_count - 1), score = upvote_count - (downvote_count - 1) WHERE uri = $1 AND deleted_at IS NULL` + decrementQuery = `UPDATE comments SET downvote_count = GREATEST(0, downvote_count - 1), score = upvote_count - GREATEST(0, downvote_count - 1) + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL` } } if decrementQuery != "" { @@ -252,14 +252,14 @@ func (c *VoteEventConsumer) indexVoteAndUpdateCounts(ctx context.Context, vote * updateQuery = ` UPDATE posts SET upvote_count = upvote_count + 1, - score = upvote_count + 1 - downvote_count + score = upvote_count + 1 - downvote_count + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } else { // "down" updateQuery = ` UPDATE posts SET downvote_count = downvote_count + 1, - score = upvote_count - (downvote_count + 1) + score = upvote_count - (downvote_count + 1) + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } @@ -270,14 +270,14 @@ func (c *VoteEventConsumer) indexVoteAndUpdateCounts(ctx context.Context, vote * updateQuery = ` UPDATE comments SET upvote_count = upvote_count + 1, - score = upvote_count + 1 - downvote_count + score = upvote_count + 1 - downvote_count + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } else { // "down" updateQuery = ` UPDATE comments SET downvote_count = downvote_count + 1, - score = upvote_count - (downvote_count + 1) + score = upvote_count - (downvote_count + 1) + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } @@ -365,14 +365,14 @@ func (c *VoteEventConsumer) deleteVoteAndUpdateCounts(ctx context.Context, vote updateQuery = ` UPDATE posts SET upvote_count = GREATEST(0, upvote_count - 1), - score = GREATEST(0, upvote_count - 1) - downvote_count + score = GREATEST(0, upvote_count - 1) - downvote_count + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } else { // "down" updateQuery = ` UPDATE posts SET downvote_count = GREATEST(0, downvote_count - 1), - score = upvote_count - GREATEST(0, downvote_count - 1) + score = upvote_count - GREATEST(0, downvote_count - 1) + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } @@ -383,14 +383,14 @@ func (c *VoteEventConsumer) deleteVoteAndUpdateCounts(ctx context.Context, vote updateQuery = ` UPDATE comments SET upvote_count = GREATEST(0, upvote_count - 1), - score = GREATEST(0, upvote_count - 1) - downvote_count + score = GREATEST(0, upvote_count - 1) - downvote_count + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } else { // "down" updateQuery = ` UPDATE comments SET downvote_count = GREATEST(0, downvote_count - 1), - score = upvote_count - GREATEST(0, downvote_count - 1) + score = upvote_count - GREATEST(0, downvote_count - 1) + bridged_upvote_count - bridged_downvote_count WHERE uri = $1 AND deleted_at IS NULL ` } diff --git a/internal/atproto/lexicon/social/coves/community/comment.json b/internal/atproto/lexicon/social/coves/community/comment.json index d774fcf..a68ba59 100644 --- a/internal/atproto/lexicon/social/coves/community/comment.json +++ b/internal/atproto/lexicon/social/coves/community/comment.json @@ -55,10 +55,37 @@ "type": "string", "format": "datetime", "description": "Timestamp of comment creation" + }, + "bridgedStats": { + "type": "ref", + "ref": "#bridgedStats", + "description": "Bridge-asserted aggregate of origin-platform votes for federated/bridged content. Set by the bridge that materialized this record; absent for natively-authored comments." } } } }, + "bridgedStats": { + "type": "object", + "description": "Aggregate vote counts asserted by the bridge for content federated from an origin platform (e.g. Lemmy). These supplement, and are kept separate from, native atproto votes.", + "required": ["upvotes", "downvotes", "asOf"], + "properties": { + "upvotes": { + "type": "integer", + "minimum": 0, + "description": "Number of upvotes on the origin platform as of asOf" + }, + "downvotes": { + "type": "integer", + "minimum": 0, + "description": "Number of downvotes on the origin platform as of asOf" + }, + "asOf": { + "type": "string", + "format": "datetime", + "description": "Timestamp the origin-platform counts were sampled; used to discard stale updates" + } + } + }, "replyRef": { "type": "object", "description": "References for maintaining thread structure. Root always points to the original post, parent points to the immediate parent (post or comment).", diff --git a/internal/atproto/lexicon/social/coves/community/post.json b/internal/atproto/lexicon/social/coves/community/post.json index 78b4947..a2041ce 100644 --- a/internal/atproto/lexicon/social/coves/community/post.json +++ b/internal/atproto/lexicon/social/coves/community/post.json @@ -92,9 +92,36 @@ "type": "string", "format": "datetime", "description": "Timestamp of post creation" + }, + "bridgedStats": { + "type": "ref", + "ref": "#bridgedStats", + "description": "Bridge-asserted aggregate of origin-platform votes for federated/bridged content. Set by the bridge that materialized this record; absent for natively-authored posts." } } } + }, + "bridgedStats": { + "type": "object", + "description": "Aggregate vote counts asserted by the bridge for content federated from an origin platform (e.g. Lemmy). These supplement, and are kept separate from, native atproto votes.", + "required": ["upvotes", "downvotes", "asOf"], + "properties": { + "upvotes": { + "type": "integer", + "minimum": 0, + "description": "Number of upvotes on the origin platform as of asOf" + }, + "downvotes": { + "type": "integer", + "minimum": 0, + "description": "Number of downvotes on the origin platform as of asOf" + }, + "asOf": { + "type": "string", + "format": "datetime", + "description": "Timestamp the origin-platform counts were sampled; used to discard stale updates" + } + } } } } diff --git a/internal/core/comments/comment.go b/internal/core/comments/comment.go index 762e41c..0c99f7f 100644 --- a/internal/core/comments/comment.go +++ b/internal/core/comments/comment.go @@ -37,20 +37,27 @@ type Comment struct { DownvoteCount int `json:"downvoteCount" db:"downvote_count"` Score int `json:"score" db:"score"` ReplyCount int `json:"replyCount" db:"reply_count"` + + // Bridge-asserted origin-platform vote aggregates for federated/bridged content. + // Populated from the record's bridgedStats field; kept separate from native votes. + // BridgedStatsAsOf is nil when no bridgedStats have ever been applied. + BridgedUpvoteCount int `json:"bridgedUpvoteCount" db:"bridged_upvote_count"` + BridgedDownvoteCount int `json:"bridgedDownvoteCount" db:"bridged_downvote_count"` + BridgedStatsAsOf *time.Time `json:"bridgedStatsAsOf,omitempty" db:"bridged_stats_as_of"` } // CommentRecord represents the atProto record structure indexed from Jetstream // This is the data structure that gets stored in the user's repository // Matches social.coves.community.comment lexicon type CommentRecord struct { - Embed interface{} `json:"embed,omitempty"` - Labels *SelfLabels `json:"labels,omitempty"` - Reply ReplyRef `json:"reply"` - Type string `json:"$type"` - Content string `json:"content"` - CreatedAt string `json:"createdAt"` + Embed interface{} `json:"embed,omitempty"` + Labels *SelfLabels `json:"labels,omitempty"` + Reply ReplyRef `json:"reply"` + Type string `json:"$type"` + Content string `json:"content"` + CreatedAt string `json:"createdAt"` Facets []interface{} `json:"facets,omitempty"` - Langs []string `json:"langs,omitempty"` + Langs []string `json:"langs,omitempty"` } // ReplyRef represents the threading structure from the comment lexicon diff --git a/internal/core/posts/post.go b/internal/core/posts/post.go index 8422c89..43fa404 100644 --- a/internal/core/posts/post.go +++ b/internal/core/posts/post.go @@ -39,6 +39,13 @@ type Post struct { DownvoteCount int `json:"downvoteCount" db:"downvote_count"` Score int `json:"score" db:"score"` CommentCount int `json:"commentCount" db:"comment_count"` + + // Bridge-asserted origin-platform vote aggregates for federated/bridged content. + // Populated from the record's bridgedStats field; kept separate from native votes. + // BridgedStatsAsOf is nil when no bridgedStats have ever been applied. + BridgedUpvoteCount int `json:"bridgedUpvoteCount" db:"bridged_upvote_count"` + BridgedDownvoteCount int `json:"bridgedDownvoteCount" db:"bridged_downvote_count"` + BridgedStatsAsOf *time.Time `json:"bridgedStatsAsOf,omitempty" db:"bridged_stats_as_of"` } // CreatePostRequest represents input for creating a new post diff --git a/internal/db/migrations/031_add_bridged_vote_stats.sql b/internal/db/migrations/031_add_bridged_vote_stats.sql new file mode 100644 index 0000000..1edfc66 --- /dev/null +++ b/internal/db/migrations/031_add_bridged_vote_stats.sql @@ -0,0 +1,57 @@ +-- +goose Up +-- Add bridge-asserted vote aggregates to posts and comments. +-- +-- Federated/bridged content (e.g. Lemmy posts materialized by the tidepool bridge) +-- carries an optional `bridgedStats` field on its record asserting the origin +-- platform's vote counts. We store those counts separately from native atproto +-- votes so the two can never clobber each other, and fold them into the displayed +-- stats and the denormalized `score` at read/write time. +-- +-- Invariant after this migration: +-- score = (upvote_count + bridged_upvote_count) - (downvote_count + bridged_downvote_count) +-- +-- bridged_stats_as_of records when the origin counts were sampled; incoming record +-- updates only overwrite the bridged counts when their asOf is newer-or-equal to the +-- stored one, which discards strictly-older/out-of-order updates while remaining +-- idempotent on replays (equal asOf carries identical counts) and tolerant of asOf +-- string-equality collisions from timestamp truncation. +-- +-- SECURITY: only records whose repo is hosted on a trusted bridge PDS may assert +-- bridgedStats (enforced in the Jetstream consumers via a provenance gate); a negative +-- or absurdly large count is rejected at the parse boundary and by the CHECK below, so +-- no native repo can self-inflate its own score. + +ALTER TABLE posts + ADD COLUMN bridged_upvote_count INT NOT NULL DEFAULT 0, + ADD COLUMN bridged_downvote_count INT NOT NULL DEFAULT 0, + ADD COLUMN bridged_stats_as_of TIMESTAMPTZ, + -- Bridged counts are unsigned aggregates; a negative value can only come from a + -- malformed/hostile record. The consumer already clamps at the parse boundary, but + -- this constraint is the last line of defence and keeps the score invariant sane. + ADD CONSTRAINT posts_bridged_counts_non_negative + CHECK (bridged_upvote_count >= 0 AND bridged_downvote_count >= 0); + +ALTER TABLE comments + ADD COLUMN bridged_upvote_count INT NOT NULL DEFAULT 0, + ADD COLUMN bridged_downvote_count INT NOT NULL DEFAULT 0, + ADD COLUMN bridged_stats_as_of TIMESTAMPTZ, + ADD CONSTRAINT comments_bridged_counts_non_negative + CHECK (bridged_upvote_count >= 0 AND bridged_downvote_count >= 0); + +COMMENT ON COLUMN posts.bridged_upvote_count IS 'Origin-platform upvotes asserted by the bridge (folded into displayed upvotes and score)'; +COMMENT ON COLUMN posts.bridged_downvote_count IS 'Origin-platform downvotes asserted by the bridge (folded into displayed downvotes and score)'; +COMMENT ON COLUMN posts.bridged_stats_as_of IS 'When the bridged counts were sampled; updates apply only when strictly newer'; +COMMENT ON COLUMN comments.bridged_upvote_count IS 'Origin-platform upvotes asserted by the bridge (folded into displayed upvotes and score)'; +COMMENT ON COLUMN comments.bridged_downvote_count IS 'Origin-platform downvotes asserted by the bridge (folded into displayed downvotes and score)'; +COMMENT ON COLUMN comments.bridged_stats_as_of IS 'When the bridged counts were sampled; updates apply only when strictly newer'; + +-- +goose Down +ALTER TABLE posts + DROP COLUMN IF EXISTS bridged_upvote_count, + DROP COLUMN IF EXISTS bridged_downvote_count, + DROP COLUMN IF EXISTS bridged_stats_as_of; + +ALTER TABLE comments + DROP COLUMN IF EXISTS bridged_upvote_count, + DROP COLUMN IF EXISTS bridged_downvote_count, + DROP COLUMN IF EXISTS bridged_stats_as_of; diff --git a/internal/db/postgres/comment_repo.go b/internal/db/postgres/comment_repo.go index 41ada92..73e2584 100644 --- a/internal/db/postgres/comment_repo.go +++ b/internal/db/postgres/comment_repo.go @@ -68,7 +68,16 @@ func (r *postgresCommentRepo) Create(ctx context.Context, comment *comments.Comm // Update modifies an existing comment's content fields // Called by Jetstream consumer after comment is updated on PDS -// Preserves vote counts and created_at timestamp +// Preserves native vote counts and created_at timestamp. +// +// The consumer passes the INCOMING bridgedStats candidate in comment.Bridged* and +// comment.BridgedStatsAsOf (nil asOf => the record carried no applicable bridgedStats). +// The bridged columns are overwritten ATOMICALLY only when an incoming asOf is present +// and is newer-or-equal to the stored bridged_stats_as_of (NULL stored => first +// application). Doing the regression comparison inside this UPDATE (rather than +// read-check-write in the consumer) makes it race-free. score is always recomputed +// inclusive of native + whichever bridged counts win, so concurrent native votes are +// never clobbered. $10 is the incoming asOf; the stored asOf is read from the row. func (r *postgresCommentRepo) Update(ctx context.Context, comment *comments.Comment) error { query := ` UPDATE comments @@ -78,7 +87,23 @@ func (r *postgresCommentRepo) Update(ctx context.Context, comment *comments.Comm content_facets = $3, embed = $4, content_labels = $5, - langs = $6 + langs = $6, + bridged_upvote_count = CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $8 ELSE bridged_upvote_count END, + bridged_downvote_count = CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $9 ELSE bridged_downvote_count END, + bridged_stats_as_of = CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $10 ELSE bridged_stats_as_of END, + score = upvote_count - downvote_count + + CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $8 ELSE bridged_upvote_count END + - CASE + WHEN $10::timestamptz IS NOT NULL AND (bridged_stats_as_of IS NULL OR $10 >= bridged_stats_as_of) + THEN $9 ELSE bridged_downvote_count END WHERE uri = $7 AND deleted_at IS NULL RETURNING id, indexed_at, created_at, upvote_count, downvote_count, score, reply_count ` @@ -92,6 +117,9 @@ func (r *postgresCommentRepo) Update(ctx context.Context, comment *comments.Comm comment.ContentLabels, pq.Array(comment.Langs), comment.URI, + comment.BridgedUpvoteCount, + comment.BridgedDownvoteCount, + comment.BridgedStatsAsOf, ).Scan( &comment.ID, &comment.IndexedAt, @@ -113,7 +141,14 @@ func (r *postgresCommentRepo) Update(ctx context.Context, comment *comments.Comm } // GetByURI retrieves a comment by its AT-URI -// Used by Jetstream consumer for UPDATE/DELETE operations +// +// This is a DISPLAY read: like every other read path (ListByParent, ListByRoot, +// GetByURIsBatch, ...) it folds the bridge-asserted aggregates into the displayed +// upvote/downvote counts (upvote_count + bridged_upvote_count, etc.) so federated +// content shows the origin platform's votes. The Jetstream update path does NOT use +// this to run its regression guard — it issues a dedicated raw-columns query so it can +// reason about the separate native and bridged aggregates (see comment_consumer's +// updateComment); folding here would conflate them. func (r *postgresCommentRepo) GetByURI(ctx context.Context, uri string) (*comments.Comment, error) { query := ` SELECT @@ -121,7 +156,7 @@ func (r *postgresCommentRepo) GetByURI(ctx context.Context, uri string) (*commen root_uri, root_cid, parent_uri, parent_cid, content, content_facets, embed, content_labels, langs, created_at, indexed_at, deleted_at, deletion_reason, deleted_by, - upvote_count, downvote_count, score, reply_count + upvote_count + bridged_upvote_count AS upvote_count, downvote_count + bridged_downvote_count AS downvote_count, score, reply_count FROM comments WHERE uri = $1 ` @@ -241,7 +276,7 @@ func (r *postgresCommentRepo) ListByRoot(ctx context.Context, rootURI string, li root_uri, root_cid, parent_uri, parent_cid, content, content_facets, embed, content_labels, langs, created_at, indexed_at, deleted_at, deletion_reason, deleted_by, - upvote_count, downvote_count, score, reply_count + upvote_count + bridged_upvote_count AS upvote_count, downvote_count + bridged_downvote_count AS downvote_count, score, reply_count FROM comments WHERE root_uri = $1 ORDER BY created_at ASC @@ -295,7 +330,7 @@ func (r *postgresCommentRepo) ListByParent(ctx context.Context, parentURI string root_uri, root_cid, parent_uri, parent_cid, content, content_facets, embed, content_labels, langs, created_at, indexed_at, deleted_at, deletion_reason, deleted_by, - upvote_count, downvote_count, score, reply_count + upvote_count + bridged_upvote_count AS upvote_count, downvote_count + bridged_downvote_count AS downvote_count, score, reply_count FROM comments WHERE parent_uri = $1 ORDER BY created_at ASC @@ -367,7 +402,7 @@ func (r *postgresCommentRepo) ListByCommenter(ctx context.Context, commenterDID root_uri, root_cid, parent_uri, parent_cid, content, content_facets, embed, content_labels, langs, created_at, indexed_at, deleted_at, deletion_reason, deleted_by, - upvote_count, downvote_count, score, reply_count + upvote_count + bridged_upvote_count AS upvote_count, downvote_count + bridged_downvote_count AS downvote_count, score, reply_count FROM comments WHERE commenter_did = $1 AND deleted_at IS NULL ORDER BY created_at DESC @@ -442,7 +477,7 @@ func (r *postgresCommentRepo) ListByCommenterWithCursor(ctx context.Context, req c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, COALESCE(u.handle, c.commenter_did) as author_handle FROM comments c LEFT JOIN users u ON c.commenter_did = u.did @@ -600,7 +635,7 @@ func (r *postgresCommentRepo) ListByParentWithHotRank( c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, log(greatest(2, c.score + 2)) / power(((EXTRACT(EPOCH FROM (NOW() - c.created_at)) / 3600) + 2), 1.8) as hot_rank, COALESCE(u.handle, c.commenter_did) as author_handle FROM comments c` @@ -611,7 +646,7 @@ func (r *postgresCommentRepo) ListByParentWithHotRank( c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, NULL::numeric as hot_rank, COALESCE(u.handle, c.commenter_did) as author_handle FROM comments c` @@ -937,7 +972,7 @@ func (r *postgresCommentRepo) GetByRootAndRkey(ctx context.Context, rootURI, rke c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, COALESCE(u.handle, c.commenter_did) as author_handle FROM comments c LEFT JOIN users u ON c.commenter_did = u.did @@ -1013,7 +1048,7 @@ func (r *postgresCommentRepo) GetByURIsBatch(ctx context.Context, uris []string) c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, COALESCE(u.handle, c.commenter_did) as author_handle FROM comments c LEFT JOIN users u ON c.commenter_did = u.did @@ -1084,7 +1119,7 @@ func (r *postgresCommentRepo) ListByParentsBatch( c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, log(greatest(2, c.score + 2)) / power(((EXTRACT(EPOCH FROM (NOW() - c.created_at)) / 3600) + 2), 1.8) as hot_rank, COALESCE(u.handle, c.commenter_did) as author_handle` // CRITICAL: Must inline hot_rank formula - PostgreSQL doesn't allow SELECT aliases in window ORDER BY @@ -1095,7 +1130,7 @@ func (r *postgresCommentRepo) ListByParentsBatch( c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, NULL::numeric as hot_rank, COALESCE(u.handle, c.commenter_did) as author_handle` windowOrderBy = `c.score DESC, c.created_at DESC` @@ -1105,7 +1140,7 @@ func (r *postgresCommentRepo) ListByParentsBatch( c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, NULL::numeric as hot_rank, COALESCE(u.handle, c.commenter_did) as author_handle` windowOrderBy = `c.created_at DESC` @@ -1116,7 +1151,7 @@ func (r *postgresCommentRepo) ListByParentsBatch( c.root_uri, c.root_cid, c.parent_uri, c.parent_cid, c.content, c.content_facets, c.embed, c.content_labels, c.langs, c.created_at, c.indexed_at, c.deleted_at, c.deletion_reason, c.deleted_by, - c.upvote_count, c.downvote_count, c.score, c.reply_count, + c.upvote_count + c.bridged_upvote_count AS upvote_count, c.downvote_count + c.bridged_downvote_count AS downvote_count, c.score, c.reply_count, log(greatest(2, c.score + 2)) / power(((EXTRACT(EPOCH FROM (NOW() - c.created_at)) / 3600) + 2), 1.8) as hot_rank, COALESCE(u.handle, c.commenter_did) as author_handle` // CRITICAL: Must inline hot_rank formula - PostgreSQL doesn't allow SELECT aliases in window ORDER BY diff --git a/internal/db/postgres/discover_repo.go b/internal/db/postgres/discover_repo.go index 37380f6..676abad 100644 --- a/internal/db/postgres/discover_repo.go +++ b/internal/db/postgres/discover_repo.go @@ -57,7 +57,7 @@ func (r *postgresDiscoverRepo) GetDiscover(ctx context.Context, req discover.Get p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count, + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count, %s as hot_rank FROM posts p`, discoverHotRankExpression) } else { @@ -68,7 +68,7 @@ func (r *postgresDiscoverRepo) GetDiscover(ctx context.Context, req discover.Get p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count, + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count, NULL::numeric as hot_rank FROM posts p` } diff --git a/internal/db/postgres/feed_repo.go b/internal/db/postgres/feed_repo.go index d0bd01d..55140c1 100644 --- a/internal/db/postgres/feed_repo.go +++ b/internal/db/postgres/feed_repo.go @@ -62,7 +62,7 @@ func (r *postgresFeedRepo) GetCommunityFeed(ctx context.Context, req communityFe p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count, + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count, %s as hot_rank FROM posts p`, communityFeedHotRankExpression) } else { @@ -73,7 +73,7 @@ func (r *postgresFeedRepo) GetCommunityFeed(ctx context.Context, req communityFe p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count, + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count, NULL::numeric as hot_rank FROM posts p` } diff --git a/internal/db/postgres/post_repo.go b/internal/db/postgres/post_repo.go index 22caa5b..f6adb3a 100644 --- a/internal/db/postgres/post_repo.go +++ b/internal/db/postgres/post_repo.go @@ -25,13 +25,18 @@ type postgresPostRepo struct { // It MUST stay byte-aligned with the positional Scan in scanPostView. It is the single // source of truth shared by GetViewsByURIs and GetByAuthor, so the column order can only // be defined in one place and the two queries cannot drift apart and silently mis-scan. +// +// Displayed upvote/downvote counts fold in the bridge-asserted aggregates +// (upvote_count + bridged_upvote_count, etc.) so federated/bridged content shows the +// origin platform's votes. score is already stored inclusive of bridged aggregates, so +// it is selected as-is. const postViewSelectColumns = ` p.uri, p.cid, p.rkey, p.author_did, u.handle as author_handle, p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count` + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count` // NewPostRepository creates a new PostgreSQL post repository func NewPostRepository(db *sql.DB) posts.Repository { @@ -112,7 +117,7 @@ func (r *postgresPostRepo) GetByURI(ctx context.Context, uri string) (*posts.Pos id, uri, cid, rkey, author_did, community_did, title, content, content_facets, embed, content_labels, created_at, edited_at, indexed_at, deleted_at, - upvote_count, downvote_count, score, comment_count + upvote_count + bridged_upvote_count AS upvote_count, downvote_count + bridged_downvote_count AS downvote_count, score, comment_count FROM posts WHERE uri = $1 ` diff --git a/internal/db/postgres/timeline_repo.go b/internal/db/postgres/timeline_repo.go index da35436..54c30b6 100644 --- a/internal/db/postgres/timeline_repo.go +++ b/internal/db/postgres/timeline_repo.go @@ -60,7 +60,7 @@ func (r *postgresTimelineRepo) GetTimeline(ctx context.Context, req timeline.Get p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count, + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count, %s as hot_rank FROM posts p`, timelineHotRankExpression) } else { @@ -71,7 +71,7 @@ func (r *postgresTimelineRepo) GetTimeline(ctx context.Context, req timeline.Get p.community_did, c.handle as community_handle, c.name as community_name, c.avatar_cid as community_avatar, c.pds_url as community_pds_url, p.title, p.content, p.content_facets, p.embed, p.content_labels, p.created_at, p.edited_at, p.indexed_at, - p.upvote_count, p.downvote_count, p.score, p.comment_count, + p.upvote_count + p.bridged_upvote_count AS upvote_count, p.downvote_count + p.bridged_downvote_count AS downvote_count, p.score, p.comment_count, NULL::numeric as hot_rank FROM posts p` } -- 2.51.2