From deff4990126d76b9b5585152a878132818a59295 Mon Sep 17 00:00:00 2001 From: Bretton <36870434+BrettM86@users.noreply.github.com> Date: Mon, 20 Jul 2026 00:12:23 -0700 Subject: [PATCH] feat(users): hydrate author profiles in all views + backfill missed profile events MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Author avatars/display names never rendered anywhere except the standalone profile endpoint, and ~870 pre-relay-switchover users had permanently bare profiles because their social.coves.actor.profile firehose event was dropped while bsky.network throttled the bridge PDSs (profile records are written once at rkey "self" — a missed create event never replays). Gap 1 — views never selected author profile columns: - postViewSelectColumns now includes u.display_name/u.avatar_cid/u.pds_url and is the single source of truth for every PostView query: the six hand-copied feed SELECT blocks collapse into feedPostSelectClause, and scanFeedPost delegates to the shared scanPostView (variadic extraDest carries hot_rank) - comment views hydrate author display name/avatar from the already batch-loaded user rows (was hardcoded nil); thread-root posts too - side effects of the unification: feeds now populate editedAt (was scanned and dropped), and malformed embed JSON omits the record key instead of serializing "embed": null Gap 2 — no reconciliation for missed profile events: - users.FetchProfileRecord fetches actor.profile/self from the user's own PDS via com.atproto.repo.getRecord (1 MiB body cap, quoted error bodies, strict RecordNotFound discrimination — bare 404s from non-PDS hosts are errors) - IndexUser backfills users indexed with a completely empty profile: sync emptiness check, detached fetch goroutine (WithoutCancel + 10s timeout) so OAuth callbacks and firehose consumers never block, pre-write re-check so a concurrent firehose event wins; opt-in via WithProfileBackfill, wired in prod - cmd/backfill-profiles: re-runnable reconciliation job for existing bare users (dry-run, bounded concurrency, atomic stats, exit 1 on failures) Hardening from multi-model review: fetched displayName/bio rune-truncated to the DB CHECK limits (hostile PDS could otherwise wedge backfill in a permanent fetch-fail loop), implausible blob CIDs rejected, HydrateImageURL no longer warn-floods on empty CIDs, social.coves.actor.profile constant consolidated to users.ProfileCollection. Tests: 20+ new unit tests (fetch classification, truncation, async backfill semantics, comment/post author hydration) plus an integration test pinning author hydration through GetViewsByURIs, GetByAuthor, and community feeds, including the author-DID-in-URL and bare-author-stays-bare cases. Also gitignores .claude/agent-memory/ (per-machine review-agent state). Deploy note: run cmd/backfill-profiles once against prod (with -dry-run first) to heal existing bare users. Co-Authored-By: Claude Fable 5 --- .gitignore | 3 + cmd/backfill-profiles/main.go | 195 ++++++++ cmd/server/main.go | 6 +- internal/api/handlers/user/update_profile.go | 7 +- internal/atproto/jetstream/feeds.go | 2 +- internal/atproto/jetstream/user_consumer.go | 6 +- internal/core/blobs/types.go | 7 + internal/core/blobs/types_test.go | 31 ++ internal/core/comments/comment_service.go | 32 +- .../core/comments/comment_service_test.go | 100 ++++ internal/core/users/profile_backfill.go | 239 +++++++++ internal/core/users/profile_backfill_test.go | 455 ++++++++++++++++++ internal/core/users/service.go | 130 ++++- internal/db/postgres/discover_repo.go | 24 +- internal/db/postgres/feed_repo.go | 24 +- internal/db/postgres/feed_repo_base.go | 137 +----- internal/db/postgres/post_repo.go | 67 ++- internal/db/postgres/timeline_repo.go | 24 +- .../author_avatar_hydration_test.go | 105 ++++ 19 files changed, 1358 insertions(+), 236 deletions(-) create mode 100644 cmd/backfill-profiles/main.go create mode 100644 internal/core/users/profile_backfill.go create mode 100644 internal/core/users/profile_backfill_test.go create mode 100644 tests/integration/author_avatar_hydration_test.go diff --git a/.gitignore b/.gitignore index 67f3592..2234f96 100644 --- a/.gitignore +++ b/.gitignore @@ -58,5 +58,8 @@ Thumbs.db # Local-only Claude Code commands (contains prod host references) .claude/commands/deploy.md +# Per-machine Claude Code agent memory (local review-agent state, not shared config) +.claude/agent-memory/ + # Beads issue tracker (local runtime state) .beads/ \ No newline at end of file diff --git a/cmd/backfill-profiles/main.go b/cmd/backfill-profiles/main.go new file mode 100644 index 0000000..b53415c --- /dev/null +++ b/cmd/backfill-profiles/main.go @@ -0,0 +1,195 @@ +// cmd/backfill-profiles/main.go +// One-off reconciliation job for users indexed without profile data. +// +// A social.coves.actor.profile record is written once at rkey "self", so its +// firehose event fires exactly once. Users whose profile event was missed +// (e.g. while bsky.network was throttling bridge PDSs) were indexed with only +// DID/handle/PDS and stay bare forever — nothing replays a lost profile commit. +// +// This job finds every user with a completely empty profile and fetches their +// profile record directly from their own PDS via com.atproto.repo.getRecord. +// It is safe to re-run: only users still lacking all profile fields are touched, +// and users without a profile record are simply skipped. +// +// If DATABASE_URL is unset, the job falls back to the local dev database +// (localhost:5435/coves_dev) and logs a prominent warning — set it explicitly +// when running against any other environment. +// +// Usage: +// +// DATABASE_URL=postgres://... go run ./cmd/backfill-profiles [-dry-run] [-concurrency N] +package main + +import ( + "context" + "database/sql" + "flag" + "log" + "net/http" + "os" + "sync" + "sync/atomic" + "time" + + "Coves/internal/core/users" + "Coves/internal/db/postgres" + + _ "github.com/lib/pq" +) + +type bareUser struct { + did string + pdsURL string +} + +// backfillStats aggregates the concurrent workers' outcome counters. +type backfillStats struct { + updated atomic.Int64 + noRecord atomic.Int64 + failed atomic.Int64 + processed atomic.Int64 +} + +func main() { + dryRun := flag.Bool("dry-run", false, "fetch and report what would be updated without writing to the database") + concurrency := flag.Int("concurrency", 4, "number of concurrent PDS fetches") + flag.Parse() + + if *concurrency < 1 { + *concurrency = 1 + } + + dbURL := os.Getenv("DATABASE_URL") + if dbURL == "" { + dbURL = "postgres://dev_user:dev_password@localhost:5435/coves_dev?sslmode=disable" + log.Printf("WARNING: DATABASE_URL is not set — falling back to the LOCAL DEV database (localhost:5435/coves_dev). Set DATABASE_URL explicitly to target any other environment.") + } + + db, err := sql.Open("postgres", dbURL) + if err != nil { + log.Fatalf("Failed to connect to database: %v", err) + } + defer func() { + if closeErr := db.Close(); closeErr != nil { + log.Printf("Warning: failed to close database: %v", closeErr) + } + }() + + ctx := context.Background() + userRepo := postgres.NewUserRepository(db) + + // Users with a completely empty profile: either their profile event was + // missed, or they genuinely have no profile record (distinguished below by + // what their PDS returns). + rows, err := db.QueryContext(ctx, ` + SELECT did, pds_url + FROM users + WHERE COALESCE(display_name, '') = '' + AND COALESCE(bio, '') = '' + AND COALESCE(avatar_cid, '') = '' + AND COALESCE(banner_cid, '') = '' + ORDER BY did + `) + if err != nil { + log.Fatalf("Failed to query users without profile data: %v", err) + } + + var candidates []bareUser + for rows.Next() { + var u bareUser + if err := rows.Scan(&u.did, &u.pdsURL); err != nil { + log.Fatalf("Failed to scan user row: %v", err) + } + candidates = append(candidates, u) + } + if err := rows.Err(); err != nil { + log.Fatalf("Failed to iterate user rows: %v", err) + } + if closeErr := rows.Close(); closeErr != nil { + log.Printf("Warning: failed to close rows: %v", closeErr) + } + + log.Printf("Found %d users without profile data (dry-run=%v, concurrency=%d)", + len(candidates), *dryRun, *concurrency) + if len(candidates) == 0 { + return + } + + client := &http.Client{Timeout: 15 * time.Second} + + var stats backfillStats + var wg sync.WaitGroup + work := make(chan bareUser) + + for i := 0; i < *concurrency; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for u := range work { + processUser(ctx, client, userRepo, u, *dryRun, &stats) + if n := stats.processed.Add(1); n%50 == 0 { + log.Printf("Progress: %d/%d processed", n, len(candidates)) + } + } + }() + } + + for _, u := range candidates { + work <- u + } + close(work) + wg.Wait() + + log.Printf("Done: %d profiles %s, %d users have no profile record, %d failed (re-run to retry failures)", + stats.updated.Load(), updatedVerb(*dryRun), stats.noRecord.Load(), stats.failed.Load()) + if stats.failed.Load() > 0 { + os.Exit(1) + } +} + +func updatedVerb(dryRun bool) string { + if dryRun { + return "would be updated" + } + return "updated" +} + +func processUser( + ctx context.Context, + client *http.Client, + userRepo users.UserRepository, + u bareUser, + dryRun bool, + stats *backfillStats, +) { + fetchCtx, cancel := context.WithTimeout(ctx, 20*time.Second) + defer cancel() + + input, err := users.FetchProfileRecord(fetchCtx, client, u.pdsURL, u.did) + if err != nil { + log.Printf("FAIL %s (%s): %v", u.did, u.pdsURL, err) + stats.failed.Add(1) + return + } + if input == nil { + stats.noRecord.Add(1) + return + } + + if dryRun { + log.Printf("WOULD UPDATE %s: displayName=%v avatar=%v banner=%v bio=%v", + u.did, input.DisplayName != nil, input.AvatarCID != nil, input.BannerCID != nil, input.Bio != nil) + stats.updated.Add(1) + return + } + + if _, err := userRepo.UpdateProfile(fetchCtx, u.did, *input); err != nil { + log.Printf("FAIL %s: profile fetched but database update failed: %v", u.did, err) + stats.failed.Add(1) + return + } + + log.Printf("UPDATED %s: displayName=%v avatar=%v banner=%v bio=%v", + u.did, input.DisplayName != nil, input.AvatarCID != nil, input.BannerCID != nil, input.Bio != nil) + stats.updated.Add(1) +} diff --git a/cmd/server/main.go b/cmd/server/main.go index 8d39998..68abab5 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -287,8 +287,12 @@ func main() { } // Initialize user repository and service early (needed for OAuth user indexing) + // Profile backfill: users indexed with no profile data (e.g. their profile + // firehose event was missed) get social.coves.actor.profile fetched from + // their PDS asynchronously, best-effort, during IndexUser. userRepo := postgresRepo.NewUserRepository(db) - userService := users.NewUserService(userRepo, identityResolver, defaultPDS, turnstileVerifier, pdsAdminPassword) + userService := users.NewUserService(userRepo, identityResolver, defaultPDS, turnstileVerifier, pdsAdminPassword, + users.WithProfileBackfill(&http.Client{Timeout: 10 * time.Second})) // Create OAuth handler for HTTP endpoints // WithUserIndexer ensures users are indexed into local database after OAuth login diff --git a/internal/api/handlers/user/update_profile.go b/internal/api/handlers/user/update_profile.go index 1ed744a..847aea3 100644 --- a/internal/api/handlers/user/update_profile.go +++ b/internal/api/handlers/user/update_profile.go @@ -10,14 +10,15 @@ import ( "Coves/internal/api/middleware" "Coves/internal/atproto/pds" + "Coves/internal/core/users" "github.com/bluesky-social/indigo/atproto/auth/oauth" ) // CovesProfileCollection is the atProto collection for Coves user profiles. -// NOTE: This constant is intentionally duplicated in internal/atproto/jetstream/user_consumer.go -// to avoid circular dependencies between packages. Keep both definitions in sync. -const CovesProfileCollection = "social.coves.actor.profile" +// NOTE: Alias of users.ProfileCollection, the canonical definition — kept as an +// exported constant of this package because existing callers reference it here. +const CovesProfileCollection = users.ProfileCollection // PDSClientFactory creates PDS clients from session data. // Used to allow injection of different auth mechanisms (OAuth for production, password for E2E tests). diff --git a/internal/atproto/jetstream/feeds.go b/internal/atproto/jetstream/feeds.go index 392c9ee..cacca95 100644 --- a/internal/atproto/jetstream/feeds.go +++ b/internal/atproto/jetstream/feeds.go @@ -54,7 +54,7 @@ var feedKeyPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]*$`) // that class of drift. var consumerWantedCollections = map[string][]string{ ConsumerUsers: { - "social.coves.actor.profile", + CovesProfileCollection, "social.coves.actor.block", }, ConsumerCommunities: { diff --git a/internal/atproto/jetstream/user_consumer.go b/internal/atproto/jetstream/user_consumer.go index 56b21df..949b3e9 100644 --- a/internal/atproto/jetstream/user_consumer.go +++ b/internal/atproto/jetstream/user_consumer.go @@ -15,9 +15,9 @@ import ( ) // CovesProfileCollection is the atProto collection for Coves user profiles. -// NOTE: This constant is intentionally duplicated in internal/api/handlers/user/update_profile.go -// to avoid circular dependencies between packages. Keep both definitions in sync. -const CovesProfileCollection = "social.coves.actor.profile" +// NOTE: Alias of users.ProfileCollection, the canonical definition — kept as an +// exported constant of this package because existing callers reference it here. +const CovesProfileCollection = users.ProfileCollection // CovesActorBlockCollection is the atProto collection for user-to-user blocks. // Records live in the blocker's repository at at://blocker_did/social.coves.actor.block/{tid} diff --git a/internal/core/blobs/types.go b/internal/core/blobs/types.go index 19dcd86..edf33cf 100644 --- a/internal/core/blobs/types.go +++ b/internal/core/blobs/types.go @@ -54,6 +54,13 @@ func HydrateImageURL(config ImageURLConfig, pdsURL, did, cid, preset string) str return HydrateBlobURL(pdsURL, did, cid) } + // Missing DID/CID is normal data (e.g. an author with no avatar), not a proxy + // configuration problem — skip the proxy attempt (and its config warning) and + // fall back to HydrateBlobURL, which returns "" for empty inputs. + if did == "" || cid == "" { + return HydrateBlobURL(pdsURL, did, cid) + } + // Determine which base URL to use baseURL := config.ProxyBaseURL if config.CDNURL != "" { diff --git a/internal/core/blobs/types_test.go b/internal/core/blobs/types_test.go index 71b9ede..3ef298f 100644 --- a/internal/core/blobs/types_test.go +++ b/internal/core/blobs/types_test.go @@ -254,6 +254,37 @@ func TestHydrateImageURL_EmptyPresetUsesDirectURL(t *testing.T) { } } +func TestHydrateImageURL_ProxyEnabledEmptyCIDReturnsEmpty(t *testing.T) { + // Missing CID (e.g. an author with no avatar) is normal data, not a config + // problem: the function must return "" without attempting proxy URL generation + // (which would emit a spurious config warning on the hot feed path). + config := ImageURLConfig{ + ProxyEnabled: true, + ProxyBaseURL: "https://coves.social", + } + pdsURL := "https://pds.example.com" + did := "did:plc:abc123" + preset := "avatar" + + if result := HydrateImageURL(config, pdsURL, did, "", preset); result != "" { + t.Errorf("HydrateImageURL with proxy enabled and empty cid = %q, want empty string", result) + } +} + +func TestHydrateImageURL_ProxyEnabledEmptyDIDReturnsEmpty(t *testing.T) { + config := ImageURLConfig{ + ProxyEnabled: true, + ProxyBaseURL: "https://coves.social", + } + pdsURL := "https://pds.example.com" + cid := "bafyreiabc123" + preset := "avatar" + + if result := HydrateImageURL(config, pdsURL, "", cid, preset); result != "" { + t.Errorf("HydrateImageURL with proxy enabled and empty did = %q, want empty string", result) + } +} + func TestImageURLConfig(t *testing.T) { // Test that ImageURLConfig holds correct fields config := ImageURLConfig{ diff --git a/internal/core/comments/comment_service.go b/internal/core/comments/comment_service.go index 5b39237..3c06133 100644 --- a/internal/core/comments/comment_service.go +++ b/internal/core/comments/comment_service.go @@ -1,6 +1,7 @@ package comments import ( + "Coves/internal/core/blobs" "Coves/internal/core/communities" "Coves/internal/core/posts" "Coves/internal/core/users" @@ -419,6 +420,23 @@ func (s *commentService) buildThreadViews( return threadViews, nil } +// hydrateAuthorProfile fills display name and avatar on an author view from an indexed +// user row. The avatar CID is transformed to an image-proxy URL against the author's +// own PDS (avatar_small preset, matching post/feed author avatars). nil user (author +// not indexed / not found) leaves the view with DID+handle only. +func hydrateAuthorProfile(view *posts.AuthorView, user *users.User) { + if user == nil { + return + } + if user.DisplayName != "" { + displayName := user.DisplayName + view.DisplayName = &displayName + } + if avatarURL := blobs.HydrateImageURL(communities.GetImageProxyConfig(), user.PDSURL, user.DID, user.AvatarCID, "avatar_small"); avatarURL != "" { + view.Avatar = &avatarURL + } +} + // buildCommentView converts a Comment entity to a CommentView with full metadata // Constructs author view, stats, and references to parent post/comment // voteStates map contains viewer's vote state for comments (from GetVoteStateForComments) @@ -440,12 +458,8 @@ func (s *commentService) buildCommentView( authorView := &posts.AuthorView{ DID: comment.CommenterDID, Handle: authorHandle, - // DisplayName, Avatar, Reputation will be populated when user profile schema is extended - // Currently User model only has DID, Handle, PDSURL fields - DisplayName: nil, - Avatar: nil, - Reputation: nil, } + hydrateAuthorProfile(authorView, usersByDID[comment.CommenterDID]) // Build aggregated statistics stats := &CommentStats{ @@ -964,7 +978,9 @@ func (s *commentService) buildPostView(ctx context.Context, post *posts.Post, vi // Build author view - fetch user to get handle (required by lexicon) // The lexicon marks authorView.handle with format:"handle", so DIDs are invalid authorHandle := post.AuthorDID // Fallback if user not found + var author *users.User if user, err := s.userRepo.GetByDID(ctx, post.AuthorDID); err == nil { + author = user authorHandle = user.Handle } else { // Log warning but don't fail the entire request @@ -974,12 +990,8 @@ func (s *commentService) buildPostView(ctx context.Context, post *posts.Post, vi authorView := &posts.AuthorView{ DID: post.AuthorDID, Handle: authorHandle, - // DisplayName, Avatar, Reputation will be populated when user profile schema is extended - // Currently User model only has DID, Handle, PDSURL fields - DisplayName: nil, - Avatar: nil, - Reputation: nil, } + hydrateAuthorProfile(authorView, author) // Build community reference - fetch community to get name and avatar (required by lexicon) // The lexicon marks communityRef.name and handle as required, so DIDs alone are insufficient diff --git a/internal/core/comments/comment_service_test.go b/internal/core/comments/comment_service_test.go index 8b375bf..06846ec 100644 --- a/internal/core/comments/comment_service_test.go +++ b/internal/core/comments/comment_service_test.go @@ -12,6 +12,7 @@ import ( "time" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) // Mock implementations for testing @@ -1482,6 +1483,105 @@ func TestCommentService_buildCommentView_BasicFields(t *testing.T) { assert.Equal(t, 0, result.Stats.ReplyCount) } +func TestCommentService_buildCommentView_HydratesAuthorProfile(t *testing.T) { + // Setup + commentRepo := newMockCommentRepo() + userRepo := newMockUserRepo() + postRepo := newMockPostRepo() + communityRepo := newMockCommunityRepo() + + postURI := "at://did:plc:post123/app.bsky.feed.post/test" + commenterDID := "did:plc:commenter123" + comment := createTestComment("at://did:plc:commenter123/comment/1", commenterDID, "commenter.test", postURI, postURI, 0) + + usersByDID := map[string]*users.User{ + commenterDID: { + DID: commenterDID, + Handle: "commenter.test", + PDSURL: "https://pds.example.com", + DisplayName: "Commenter Display", + AvatarCID: "bafkreicommenteravatar", + }, + } + + service := NewCommentService(commentRepo, userRepo, postRepo, communityRepo, nil, nil, nil).(*commentService) + + // Execute + result := service.buildCommentView(comment, nil, nil, usersByDID) + + // Verify author profile hydration + require.NotNil(t, result.Author) + require.NotNil(t, result.Author.DisplayName, "display name must be hydrated from indexed user") + assert.Equal(t, "Commenter Display", *result.Author.DisplayName) + require.NotNil(t, result.Author.Avatar, "avatar must be hydrated from indexed user") + assert.Contains(t, *result.Author.Avatar, "bafkreicommenteravatar", "avatar URL must reference the avatar CID") + assert.Contains(t, *result.Author.Avatar, "did%3Aplc%3Acommenter123", "avatar URL must reference the author's own DID") +} + +func TestCommentService_buildPostView_HydratesAuthorProfile(t *testing.T) { + // Setup + commentRepo := newMockCommentRepo() + userRepo := newMockUserRepo() + postRepo := newMockPostRepo() + communityRepo := newMockCommunityRepo() + + authorDID := "did:plc:postauthor123" + communityDID := "did:plc:community123" + post := createTestPost("at://did:plc:postauthor123/social.coves.community.post/test", authorDID, communityDID) + + author := createTestUser(authorDID, "postauthor.test") + author.DisplayName = "Post Author Display" + author.AvatarCID = "bafkreipostauthoravatar" + author.PDSURL = "https://pds.example.com" + _, _ = userRepo.Create(context.Background(), author) + + community := createTestCommunity(communityDID, "c-test.coves.social") + _, _ = communityRepo.Create(context.Background(), community) + + service := NewCommentService(commentRepo, userRepo, postRepo, communityRepo, nil, nil, nil).(*commentService) + + // Execute + result := service.buildPostView(context.Background(), post, nil) + + // Verify author profile hydration + require.NotNil(t, result.Author) + require.NotNil(t, result.Author.DisplayName, "display name must be hydrated from indexed user") + assert.Equal(t, "Post Author Display", *result.Author.DisplayName) + require.NotNil(t, result.Author.Avatar, "avatar must be hydrated from indexed user") + assert.Contains(t, *result.Author.Avatar, "bafkreipostauthoravatar", "avatar URL must reference the avatar CID") + assert.Contains(t, *result.Author.Avatar, "did%3Aplc%3Apostauthor123", "avatar URL must reference the author's own DID") +} + +func TestCommentService_buildCommentView_NoProfileLeavesAuthorBare(t *testing.T) { + // Setup: user indexed without profile data (e.g. profile event missed) + commentRepo := newMockCommentRepo() + userRepo := newMockUserRepo() + postRepo := newMockPostRepo() + communityRepo := newMockCommunityRepo() + + postURI := "at://did:plc:post123/app.bsky.feed.post/test" + commenterDID := "did:plc:commenter123" + comment := createTestComment("at://did:plc:commenter123/comment/1", commenterDID, "commenter.test", postURI, postURI, 0) + + usersByDID := map[string]*users.User{ + commenterDID: { + DID: commenterDID, + Handle: "commenter.test", + PDSURL: "https://pds.example.com", + }, + } + + service := NewCommentService(commentRepo, userRepo, postRepo, communityRepo, nil, nil, nil).(*commentService) + + // Execute + result := service.buildCommentView(comment, nil, nil, usersByDID) + + // Verify: no fabricated profile data + require.NotNil(t, result.Author) + assert.Nil(t, result.Author.DisplayName) + assert.Nil(t, result.Author.Avatar) +} + func TestCommentService_buildCommentView_TopLevelComment(t *testing.T) { // Setup commentRepo := newMockCommentRepo() diff --git a/internal/core/users/profile_backfill.go b/internal/core/users/profile_backfill.go new file mode 100644 index 0000000..481fee9 --- /dev/null +++ b/internal/core/users/profile_backfill.go @@ -0,0 +1,239 @@ +package users + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "strconv" + "strings" +) + +// ProfileCollection is the atProto collection holding Coves user profiles. +// The profile record always lives at rkey "self" in the user's own repository. +const ProfileCollection = "social.coves.actor.profile" + +// profileRecordRKey is the fixed rkey for profile records (one profile per repo). +const profileRecordRKey = "self" + +// maxProfileResponseBytes bounds how much of a PDS getRecord response we read. +// A legitimate profile record is well under this; the cap protects against a +// misbehaving or malicious PDS streaming an unbounded body. +const maxProfileResponseBytes = 1 << 20 // 1 MiB + +// Field caps mirror the users table CHECK constraints (migration 027: +// display_name <= 64, bio <= 256 — Postgres length() counts characters, so +// these are rune counts, not bytes). Overlong values from an untrusted PDS are +// truncated rather than rejected: a truncated profile indexes fine, while an +// over-limit UPDATE would fail the constraint forever on every re-index. +const ( + maxDisplayNameChars = 64 + maxBioChars = 256 +) + +// maxBlobCIDLength caps the accepted length of a blob ref's $link. Real CIDs +// are ~59 chars (CIDv1 base32); anything past this is not a plausible CID. +const maxBlobCIDLength = 256 + +// FetchProfileRecord fetches a user's social.coves.actor.profile/self record directly +// from their PDS via com.atproto.repo.getRecord and converts it to an UpdateProfileInput. +// +// This is the reconciliation path for profile events the firehose never delivered +// (a profile record is written once at rkey "self", so its create event fires exactly +// once — if it was missed, nothing ever replays it). Returns (nil, nil) when the user +// has no profile record or the record carries no profile fields — "nothing to apply" +// is a normal outcome, not an error. +func FetchProfileRecord(ctx context.Context, client *http.Client, pdsURL, did string) (*UpdateProfileInput, error) { + if client == nil { + return nil, fmt.Errorf("http client is required") + } + if strings.TrimSpace(pdsURL) == "" { + return nil, fmt.Errorf("PDS URL is required") + } + if strings.TrimSpace(did) == "" { + return nil, fmt.Errorf("DID is required") + } + + endpoint := strings.TrimSuffix(pdsURL, "/") + "/xrpc/com.atproto.repo.getRecord?repo=" + + url.QueryEscape(did) + "&collection=" + ProfileCollection + "&rkey=" + profileRecordRKey + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil) + if err != nil { + return nil, fmt.Errorf("failed to create getRecord request: %w", err) + } + + resp, err := client.Do(req) + if err != nil { + return nil, fmt.Errorf("failed to fetch profile record from PDS: %w", err) + } + defer func() { + _ = resp.Body.Close() + }() + + body, err := io.ReadAll(io.LimitReader(resp.Body, maxProfileResponseBytes)) + if err != nil { + return nil, fmt.Errorf("failed to read getRecord response: %w", err) + } + + switch { + case resp.StatusCode == http.StatusOK: + // fall through to parse below + case resp.StatusCode == http.StatusNotFound && isXRPCErrorBody(body): + // Some PDS implementations 404 on missing records. Only trust the 404 + // when the body is an XRPC error object — a bare 404 (HTML from a + // reverse proxy or generic web server behind a stale pds_url) means we + // never reached a PDS at all, so it must surface as an error, not be + // silently classified as "user has no profile record". + return nil, nil + case resp.StatusCode == http.StatusBadRequest && isRecordNotFoundBody(body): + // Reference PDS returns 400 RecordNotFound when the repo exists but has no profile + return nil, nil + default: + // SECURITY: cap the echoed body so a hostile PDS can't flood our logs, + // and quote it so control chars / ANSI escapes can't corrupt log output. + detail := string(body) + if len(detail) > 256 { + detail = detail[:256] + } + return nil, fmt.Errorf("PDS getRecord returned status %d: %s", resp.StatusCode, strconv.Quote(detail)) + } + + var parsed struct { + Value map[string]interface{} `json:"value"` + } + if err := json.Unmarshal(body, &parsed); err != nil { + return nil, fmt.Errorf("failed to parse getRecord response: %w", err) + } + if parsed.Value == nil { + return nil, fmt.Errorf("getRecord response missing record value") + } + + input := parseProfileRecord(parsed.Value) + if input.DisplayName == nil && input.Bio == nil && input.AvatarCID == nil && input.BannerCID == nil { + return nil, nil // record exists but carries nothing we index + } + return &input, nil +} + +// isXRPCErrorBody reports whether body is an XRPC error JSON object (has a +// non-empty "error" field). Used to distinguish a real PDS not-found response +// from a bare 404 served by whatever non-PDS host a stale pds_url points at. +func isXRPCErrorBody(body []byte) bool { + var xrpcErr struct { + Error string `json:"error"` + } + if err := json.Unmarshal(body, &xrpcErr); err != nil { + return false + } + return xrpcErr.Error != "" +} + +// isRecordNotFoundBody reports whether an XRPC error body is the reference PDS's +// "record not found" rejection (as opposed to some other 400-class failure). +// The message-substring match is only trusted on the error codes the reference +// PDS actually uses for missing records (RecordNotFound, InvalidRequest) — an +// arbitrary error that merely mentions "could not locate record" is not proof +// the record is missing. +func isRecordNotFoundBody(body []byte) bool { + var xrpcErr struct { + Error string `json:"error"` + Message string `json:"message"` + } + if err := json.Unmarshal(body, &xrpcErr); err != nil { + return false + } + if xrpcErr.Error == "RecordNotFound" { + return true + } + if xrpcErr.Error != "InvalidRequest" { + return false + } + return strings.Contains(strings.ToLower(xrpcErr.Message), "could not locate record") +} + +// parseProfileRecord extracts indexable profile fields from a social.coves.actor.profile +// record value. Field extraction mirrors the Jetstream profile consumer +// (internal/atproto/jetstream/user_consumer.go handleProfileUpdate): only fields present +// in the record are set, and avatar/banner must be well-formed blob refs. +// SECURITY: the record comes from an untrusted PDS — displayName/description are +// truncated to the DB CHECK constraint limits so an overlong value indexes +// (truncated) instead of failing the UPDATE forever on every re-index. +func parseProfileRecord(record map[string]interface{}) UpdateProfileInput { + input := UpdateProfileInput{} + + if displayName, ok := record["displayName"].(string); ok { + displayName = truncateRunes(displayName, maxDisplayNameChars) + input.DisplayName = &displayName + } + if description, ok := record["description"].(string); ok { + description = truncateRunes(description, maxBioChars) + input.Bio = &description + } + if avatar, ok := record["avatar"].(map[string]interface{}); ok { + if cid, ok := extractProfileBlobCID(avatar); ok { + input.AvatarCID = &cid + } + } + if banner, ok := record["banner"].(map[string]interface{}); ok { + if cid, ok := extractProfileBlobCID(banner); ok { + input.BannerCID = &cid + } + } + + return input +} + +// extractProfileBlobCID pulls the CID out of an atProto blob ref +// ({"$type":"blob","ref":{"$link":""},...}). Returns false for anything +// that is not a well-formed blob ref or whose $link is not a plausible CID. +func extractProfileBlobCID(blob map[string]interface{}) (string, bool) { + blobType, ok := blob["$type"].(string) + if !ok || blobType != "blob" { + return "", false + } + ref, ok := blob["ref"].(map[string]interface{}) + if !ok { + return "", false + } + link, ok := ref["$link"].(string) + if !ok || !isPlausibleCID(link) { + return "", false + } + return link, true +} + +// isPlausibleCID reports whether link looks like a CID: non-empty, bounded +// length, and restricted to the alphanumeric charset of base32/base58 CID +// encodings. This is a sanity gate on untrusted PDS output, not full CID +// validation — it keeps arbitrary strings out of the avatar/banner columns. +func isPlausibleCID(link string) bool { + if link == "" || len(link) > maxBlobCIDLength { + return false + } + for _, c := range link { + isAlphanumeric := (c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c <= '9') + if !isAlphanumeric { + return false + } + } + return true +} + +// truncateRunes returns s truncated to at most maxRunes characters without +// splitting a multi-byte rune (the DB constraints count characters, so the cut +// must be by runes, not bytes). +func truncateRunes(s string, maxRunes int) string { + if len(s) <= maxRunes { + return s // byte length within the cap implies rune count is too + } + count := 0 + for i := range s { + if count == maxRunes { + return s[:i] + } + count++ + } + return s +} diff --git a/internal/core/users/profile_backfill_test.go b/internal/core/users/profile_backfill_test.go new file mode 100644 index 0000000..5a8c2f2 --- /dev/null +++ b/internal/core/users/profile_backfill_test.go @@ -0,0 +1,455 @@ +package users + +import ( + "context" + "encoding/json" + "net/http" + "net/http/httptest" + "strings" + "sync/atomic" + "testing" + "time" + "unicode/utf8" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" +) + +// backfillSettle is how long the skip-path tests wait before asserting that no +// backfill fetch happened — long enough to catch a regressed async fetch. +const backfillSettle = 100 * time.Millisecond + +// waitForBackfill waits for the async backfill goroutine to hit the fake PDS +// and (via done) finish its UpdateProfile write. +func waitForBackfill(t *testing.T, hits *int64, done <-chan struct{}) { + t.Helper() + require.Eventually(t, func() bool { + return atomic.LoadInt64(hits) == 1 + }, 5*time.Second, 10*time.Millisecond, "backfill goroutine never fetched from the PDS") + select { + case <-done: + case <-time.After(5 * time.Second): + t.Fatal("backfill goroutine never called UpdateProfile") + } +} + +// newProfilePDS spins up a fake PDS that answers com.atproto.repo.getRecord for +// profile records. handler receives the request after basic endpoint assertions. +func newProfilePDS(t *testing.T, expectedDID string, hits *int64, handler http.HandlerFunc) *httptest.Server { + t.Helper() + return httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + atomic.AddInt64(hits, 1) + assert.Equal(t, "/xrpc/com.atproto.repo.getRecord", r.URL.Path) + assert.Equal(t, expectedDID, r.URL.Query().Get("repo")) + assert.Equal(t, ProfileCollection, r.URL.Query().Get("collection")) + assert.Equal(t, "self", r.URL.Query().Get("rkey")) + handler(w, r) + })) +} + +func writeProfileRecord(t *testing.T, w http.ResponseWriter, value map[string]interface{}) { + t.Helper() + w.Header().Set("Content-Type", "application/json") + err := json.NewEncoder(w).Encode(map[string]interface{}{ + "uri": "at://did:plc:test/social.coves.actor.profile/self", + "cid": "bafyreicid", + "value": value, + }) + require.NoError(t, err) +} + +func blobRef(cid string) map[string]interface{} { + return map[string]interface{}{ + "$type": "blob", + "ref": map[string]interface{}{"$link": cid}, + "mimeType": "image/png", + "size": 12345, + } +} + +func TestFetchProfileRecord_AllFields(t *testing.T) { + var hits int64 + srv := newProfilePDS(t, "did:plc:alice", &hits, func(w http.ResponseWriter, r *http.Request) { + writeProfileRecord(t, w, map[string]interface{}{ + "$type": ProfileCollection, + "displayName": "Alice", + "description": "hello from alice", + "avatar": blobRef("bafkreiavatar"), + "banner": blobRef("bafkreibanner"), + }) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:alice") + require.NoError(t, err) + require.NotNil(t, input) + + require.NotNil(t, input.DisplayName) + assert.Equal(t, "Alice", *input.DisplayName) + require.NotNil(t, input.Bio) + assert.Equal(t, "hello from alice", *input.Bio) + require.NotNil(t, input.AvatarCID) + assert.Equal(t, "bafkreiavatar", *input.AvatarCID) + require.NotNil(t, input.BannerCID) + assert.Equal(t, "bafkreibanner", *input.BannerCID) + assert.Equal(t, int64(1), atomic.LoadInt64(&hits)) +} + +func TestFetchProfileRecord_Bare404IsError(t *testing.T) { + // A 404 whose body is NOT an XRPC error object means we likely never reached + // a PDS (stale pds_url behind a generic web server / reverse proxy). That + // must surface as an error — silently classifying it as "user has no profile + // record" would give up permanently and invisibly. + var hits int64 + srv := newProfilePDS(t, "did:plc:stale", &hits, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/html") + w.WriteHeader(http.StatusNotFound) + _, _ = w.Write([]byte("404")) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:stale") + require.Error(t, err, "a bare 404 (no XRPC error body) must be an error, not not-found") + assert.Nil(t, input) + assert.Contains(t, err.Error(), "404") +} + +func TestFetchProfileRecord_XRPC404IsNotFound(t *testing.T) { + // A 404 carrying a real XRPC error body came from a PDS — that IS a missing + // record (some PDS implementations 404 instead of 400 RecordNotFound). + var hits int64 + srv := newProfilePDS(t, "did:plc:notfound404", &hits, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusNotFound) + _, _ = w.Write([]byte(`{"error":"NotFound","message":"record not found"}`)) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:notfound404") + require.NoError(t, err) + assert.Nil(t, input, "XRPC 404 means the user has no profile record") +} + +func TestFetchProfileRecord_RecordNotFound(t *testing.T) { + var hits int64 + srv := newProfilePDS(t, "did:plc:bob", &hits, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _, _ = w.Write([]byte(`{"error":"RecordNotFound","message":"Record not found"}`)) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:bob") + require.NoError(t, err) + assert.Nil(t, input, "missing profile record must be a nil result, not an error") +} + +func TestFetchProfileRecord_CouldNotLocateRecordMessage(t *testing.T) { + // Some PDS implementations use InvalidRequest with a descriptive message + var hits int64 + srv := newProfilePDS(t, "did:plc:bob", &hits, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _, _ = w.Write([]byte(`{"error":"InvalidRequest","message":"Could not locate record: at://did:plc:bob/social.coves.actor.profile/self"}`)) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:bob") + require.NoError(t, err) + assert.Nil(t, input) +} + +func TestFetchProfileRecord_ServerErrorPropagates(t *testing.T) { + var hits int64 + srv := newProfilePDS(t, "did:plc:carol", &hits, func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + _, _ = w.Write([]byte("boom")) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:carol") + require.Error(t, err) + assert.Nil(t, input) + assert.Contains(t, err.Error(), "500") +} + +func TestFetchProfileRecord_OtherBadRequestIsError(t *testing.T) { + // A 400 that is NOT record-not-found (e.g. malformed repo) must surface as an + // error so callers don't mistake a broken request for "user has no profile" + var hits int64 + srv := newProfilePDS(t, "did:plc:dave", &hits, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _, _ = w.Write([]byte(`{"error":"InvalidRequest","message":"Invalid repo identifier"}`)) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:dave") + require.Error(t, err) + assert.Nil(t, input) +} + +func TestFetchProfileRecord_LocateMessageWithWrongErrorCodeIsError(t *testing.T) { + // The "could not locate record" message substring is only trusted on the + // error codes the reference PDS uses for missing records (RecordNotFound, + // InvalidRequest). An arbitrary error that merely mentions the phrase must + // not be classified as not-found. + var hits int64 + srv := newProfilePDS(t, "did:plc:grace", &hits, func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusBadRequest) + _, _ = w.Write([]byte(`{"error":"SomethingElse","message":"could not locate record"}`)) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:grace") + require.Error(t, err, "message substring alone must not classify a 400 as not-found") + assert.Nil(t, input) +} + +func TestParseProfileRecord_TruncatesOverlongFields(t *testing.T) { + // The users table CHECK constraints cap display_name at 64 and bio at 256 + // characters (Postgres length() counts characters, so truncation must be by + // runes). Multi-byte runes prove the cut is not by bytes and never splits a rune. + longName := strings.Repeat("é", 70) // 70 runes, 140 bytes + longDescription := strings.Repeat("日", 300) // 300 runes, 900 bytes + + input := parseProfileRecord(map[string]interface{}{ + "displayName": longName, + "description": longDescription, + }) + + require.NotNil(t, input.DisplayName) + assert.Equal(t, 64, utf8.RuneCountInString(*input.DisplayName)) + assert.Equal(t, strings.Repeat("é", 64), *input.DisplayName) + assert.True(t, utf8.ValidString(*input.DisplayName), "truncation must not split a rune") + + require.NotNil(t, input.Bio) + assert.Equal(t, 256, utf8.RuneCountInString(*input.Bio)) + assert.Equal(t, strings.Repeat("日", 256), *input.Bio) + assert.True(t, utf8.ValidString(*input.Bio), "truncation must not split a rune") + + // At-limit values pass through untouched + atLimit := parseProfileRecord(map[string]interface{}{ + "displayName": strings.Repeat("a", 64), + "description": strings.Repeat("b", 256), + }) + require.NotNil(t, atLimit.DisplayName) + assert.Equal(t, strings.Repeat("a", 64), *atLimit.DisplayName) + require.NotNil(t, atLimit.Bio) + assert.Equal(t, strings.Repeat("b", 256), *atLimit.Bio) +} + +func TestParseProfileRecord_RejectsImplausibleBlobCID(t *testing.T) { + input := parseProfileRecord(map[string]interface{}{ + "avatar": blobRef("../../../etc/passwd"), // non-alphanumeric chars + "banner": blobRef(strings.Repeat("a", 300)), // over the length cap + }) + assert.Nil(t, input.AvatarCID, "$link with non-CID charset must be rejected") + assert.Nil(t, input.BannerCID, "$link over the length cap must be rejected") + + valid := parseProfileRecord(map[string]interface{}{ + "avatar": blobRef("bafkreib2qya3v6fyfvdkr5gkuwrmhxjkkvsyyx2xczpkqnvkq3rc5jm2gq"), + }) + require.NotNil(t, valid.AvatarCID) + assert.Equal(t, "bafkreib2qya3v6fyfvdkr5gkuwrmhxjkkvsyyx2xczpkqnvkq3rc5jm2gq", *valid.AvatarCID) +} + +func TestFetchProfileRecord_MalformedBlobIgnored(t *testing.T) { + var hits int64 + srv := newProfilePDS(t, "did:plc:eve", &hits, func(w http.ResponseWriter, r *http.Request) { + writeProfileRecord(t, w, map[string]interface{}{ + "displayName": "Eve", + // Not a blob ref: missing $type/ref structure + "avatar": map[string]interface{}{"cid": "bafkreibad"}, + }) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:eve") + require.NoError(t, err) + require.NotNil(t, input) + require.NotNil(t, input.DisplayName) + assert.Equal(t, "Eve", *input.DisplayName) + assert.Nil(t, input.AvatarCID, "malformed blob ref must not produce an avatar CID") +} + +func TestFetchProfileRecord_EmptyRecordIsNil(t *testing.T) { + var hits int64 + srv := newProfilePDS(t, "did:plc:frank", &hits, func(w http.ResponseWriter, r *http.Request) { + writeProfileRecord(t, w, map[string]interface{}{"$type": ProfileCollection}) + }) + defer srv.Close() + + input, err := FetchProfileRecord(context.Background(), srv.Client(), srv.URL, "did:plc:frank") + require.NoError(t, err) + assert.Nil(t, input, "record with no indexable fields means nothing to apply") +} + +// TestIndexUser_BackfillsEmptyProfile verifies the full IndexUser → fetch → UpdateProfile +// path for a newly indexed user with no profile data. +func TestIndexUser_BackfillsEmptyProfile(t *testing.T) { + testDID := "did:plc:backfillme" + testHandle := "backfillme.test" + + var hits int64 + srv := newProfilePDS(t, testDID, &hits, func(w http.ResponseWriter, r *http.Request) { + writeProfileRecord(t, w, map[string]interface{}{ + "displayName": "Backfilled", + "avatar": blobRef("bafkreiavatar"), + }) + }) + defer srv.Close() + + mockRepo := new(MockUserRepository) + newUser := &User{ + DID: testDID, + Handle: testHandle, + PDSURL: srv.URL, // profile is fetched from the user's own PDS + CreatedAt: time.Now(), + } + mockRepo.On("Create", mock.Anything, mock.Anything).Return(newUser, nil) + // The detached goroutine re-checks emptiness before writing (a concurrent + // firehose event may have won the race) — still-empty user lets the write proceed. + mockRepo.On("GetByDID", mock.Anything, testDID).Return(newUser, nil) + done := make(chan struct{}) + mockRepo.On("UpdateProfile", mock.Anything, testDID, mock.MatchedBy(func(input UpdateProfileInput) bool { + return input.DisplayName != nil && *input.DisplayName == "Backfilled" && + input.AvatarCID != nil && *input.AvatarCID == "bafkreiavatar" && + input.Bio == nil && input.BannerCID == nil + })).Return(newUser, nil).Run(func(args mock.Arguments) { close(done) }) + + service := NewUserService(mockRepo, nil, "https://default.pds", nil, "", + WithProfileBackfill(srv.Client())) + + err := service.IndexUser(context.Background(), testDID, testHandle, srv.URL) + require.NoError(t, err) + + waitForBackfill(t, &hits, done) + mockRepo.AssertExpectations(t) +} + +// TestIndexUser_SkipsBackfillWhenProfilePopulated verifies that users who already have +// profile data are never re-fetched (the firehose is the source of truth for them). +func TestIndexUser_SkipsBackfillWhenProfilePopulated(t *testing.T) { + testDID := "did:plc:hasprofile" + + var hits int64 + srv := newProfilePDS(t, testDID, &hits, func(w http.ResponseWriter, r *http.Request) { + t.Error("PDS must not be called for a user with existing profile data") + }) + defer srv.Close() + + mockRepo := new(MockUserRepository) + existingUser := &User{ + DID: testDID, + Handle: "hasprofile.test", + PDSURL: srv.URL, + AvatarCID: "bafkreialready", + CreatedAt: time.Now(), + } + mockRepo.On("Create", mock.Anything, mock.Anything).Return(existingUser, nil) + + service := NewUserService(mockRepo, nil, "https://default.pds", nil, "", + WithProfileBackfill(srv.Client())) + + err := service.IndexUser(context.Background(), testDID, "hasprofile.test", srv.URL) + require.NoError(t, err) + + // The emptiness check is synchronous (no goroutine spawns for populated + // profiles), but settle briefly so a regressed async fetch would be caught. + time.Sleep(backfillSettle) + assert.Equal(t, int64(0), atomic.LoadInt64(&hits)) + mockRepo.AssertNotCalled(t, "UpdateProfile", mock.Anything, mock.Anything, mock.Anything) +} + +// TestIndexUser_BackfillDisabledByDefault verifies that without WithProfileBackfill, +// IndexUser never reaches out to a PDS (preserves prior behavior for all callers that +// don't opt in). +func TestIndexUser_BackfillDisabledByDefault(t *testing.T) { + testDID := "did:plc:nobackfill" + + var hits int64 + srv := newProfilePDS(t, testDID, &hits, func(w http.ResponseWriter, r *http.Request) { + t.Error("PDS must not be called when backfill is not enabled") + }) + defer srv.Close() + + mockRepo := new(MockUserRepository) + newUser := &User{DID: testDID, Handle: "nobackfill.test", PDSURL: srv.URL, CreatedAt: time.Now()} + mockRepo.On("Create", mock.Anything, mock.Anything).Return(newUser, nil) + + service := NewUserService(mockRepo, nil, "https://default.pds", nil, "") + + err := service.IndexUser(context.Background(), testDID, "nobackfill.test", srv.URL) + require.NoError(t, err) + + // Backfill disabled → no goroutine spawns; settle briefly to catch a + // regressed async fetch before asserting. + time.Sleep(backfillSettle) + assert.Equal(t, int64(0), atomic.LoadInt64(&hits)) +} + +// TestIndexUser_BackfillFailureDoesNotFailIndexing verifies backfill is best-effort: +// a dead or erroring PDS must not fail the IndexUser call that triggered it (which +// would dead-letter the post/comment event being consumed). +func TestIndexUser_BackfillFailureDoesNotFailIndexing(t *testing.T) { + testDID := "did:plc:deadpds" + + var hits int64 + srv := newProfilePDS(t, testDID, &hits, func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + }) + defer srv.Close() + + mockRepo := new(MockUserRepository) + newUser := &User{DID: testDID, Handle: "deadpds.test", PDSURL: srv.URL, CreatedAt: time.Now()} + mockRepo.On("Create", mock.Anything, mock.Anything).Return(newUser, nil) + + service := NewUserService(mockRepo, nil, "https://default.pds", nil, "", + WithProfileBackfill(srv.Client())) + + err := service.IndexUser(context.Background(), testDID, "deadpds.test", srv.URL) + assert.NoError(t, err, "backfill failure must never fail indexing") + + require.Eventually(t, func() bool { + return atomic.LoadInt64(&hits) == 1 + }, 5*time.Second, 10*time.Millisecond, "backfill goroutine never fetched from the PDS") + // Settle so the goroutine's post-fetch path (which must bail on the error) + // has finished before asserting no write happened. + time.Sleep(backfillSettle) + mockRepo.AssertNotCalled(t, "UpdateProfile", mock.Anything, mock.Anything, mock.Anything) +} + +// TestIndexUser_BackfillsOnHandleChange verifies the handle-conflict path (existing +// user, changed handle) also heals an empty profile. +func TestIndexUser_BackfillsOnHandleChange(t *testing.T) { + testDID := "did:plc:renamed" + newHandle := "newname.test" + + var hits int64 + srv := newProfilePDS(t, testDID, &hits, func(w http.ResponseWriter, r *http.Request) { + writeProfileRecord(t, w, map[string]interface{}{"displayName": "Renamed"}) + }) + defer srv.Close() + + mockRepo := new(MockUserRepository) + renamedUser := &User{DID: testDID, Handle: newHandle, PDSURL: srv.URL, CreatedAt: time.Now()} + mockRepo.On("Create", mock.Anything, mock.Anything).Return(nil, ErrHandleAlreadyTaken) + mockRepo.On("UpdateHandle", mock.Anything, testDID, newHandle).Return(renamedUser, nil) + mockRepo.On("GetByDID", mock.Anything, testDID).Return(renamedUser, nil) + done := make(chan struct{}) + mockRepo.On("UpdateProfile", mock.Anything, testDID, mock.MatchedBy(func(input UpdateProfileInput) bool { + return input.DisplayName != nil && *input.DisplayName == "Renamed" + })).Return(renamedUser, nil).Run(func(args mock.Arguments) { close(done) }) + + service := NewUserService(mockRepo, nil, "https://default.pds", nil, "", + WithProfileBackfill(srv.Client())) + + err := service.IndexUser(context.Background(), testDID, newHandle, srv.URL) + require.NoError(t, err) + + waitForBackfill(t, &hits, done) + mockRepo.AssertExpectations(t) +} diff --git a/internal/core/users/service.go b/internal/core/users/service.go index 6bf539f..c510d2f 100644 --- a/internal/core/users/service.go +++ b/internal/core/users/service.go @@ -49,6 +49,11 @@ const ( // can't pin a captcha-verified caller for the request's full lifetime. const pdsAdminCallTimeout = 10 * time.Second +// profileBackfillTimeout bounds the detached profile-backfill fetch+store. The +// backfill goroutine is decoupled from the caller's request context, so this +// is its only deadline. +const profileBackfillTimeout = 10 * time.Second + type userService struct { userRepo UserRepository identityResolver identity.Resolver @@ -66,6 +71,30 @@ type userService struct { // pdsAdminClient is reused across calls so HTTP/1.1 keep-alive and TLS // session resumption actually kick in. pdsAdminClient *http.Client + + // profileBackfillClient, when non-nil, enables profile backfill on IndexUser: + // a user indexed with no profile data gets their social.coves.actor.profile + // record fetched directly from their PDS. This reconciles profile firehose + // events that were missed (a profile is written once at rkey "self", so a + // missed create event is never re-emitted). nil → backfill disabled. + profileBackfillClient *http.Client +} + +// UserServiceOption configures optional behavior on the user service. +type UserServiceOption func(*userService) + +// WithProfileBackfill enables best-effort profile backfill during IndexUser (see +// profileBackfillClient). Pass nil to use a default client with a 10s timeout. +// The fetch+store runs in a detached goroutine so it never blocks IndexUser +// callers (OAuth login, Jetstream consumers); failures are logged only — run +// cmd/backfill-profiles to reconcile users whose backfill fetch failed. +func WithProfileBackfill(client *http.Client) UserServiceOption { + return func(s *userService) { + if client == nil { + client = &http.Client{Timeout: 10 * time.Second} + } + s.profileBackfillClient = client + } } // NewUserService creates a new user service. @@ -79,8 +108,9 @@ func NewUserService( defaultPDS string, turnstile TurnstileVerifier, pdsAdminPassword string, + opts ...UserServiceOption, ) UserService { - return &userService{ + s := &userService{ userRepo: userRepo, identityResolver: identityResolver, defaultPDS: defaultPDS, @@ -88,6 +118,10 @@ func NewUserService( pdsAdminPassword: pdsAdminPassword, pdsAdminClient: &http.Client{Timeout: pdsAdminCallTimeout}, } + for _, opt := range opts { + opt(s) + } + return s } // CreateUser creates a new user in the AppView database @@ -356,9 +390,13 @@ func (s *userService) mintInviteCode(ctx context.Context) (string, error) { // IndexUser creates or updates a user in the local database. // This is idempotent and safe to call multiple times for the same user. // If the user exists, their handle is updated if it changed. +// When profile backfill is enabled (WithProfileBackfill) and the indexed user has no +// profile data, their profile record is fetched from their PDS asynchronously and +// best-effort — this heals users whose profile firehose event was never delivered +// without ever blocking or failing the IndexUser call itself. func (s *userService) IndexUser(ctx context.Context, did, handle, pdsURL string) error { // Try to create the user (idempotent - CreateUser returns existing user if DID exists) - _, err := s.CreateUser(ctx, CreateUserRequest{ + user, err := s.CreateUser(ctx, CreateUserRequest{ DID: did, Handle: handle, PDSURL: pdsURL, @@ -367,20 +405,92 @@ func (s *userService) IndexUser(ctx context.Context, did, handle, pdsURL string) if err != nil { // Check if it's a handle conflict (user exists with different handle) // In this case, update the handle instead - if errors.Is(err, ErrHandleAlreadyTaken) { - // User exists but handle changed - update it - _, updateErr := s.UpdateHandle(ctx, did, handle) - if updateErr != nil { - return fmt.Errorf("failed to update handle for existing user: %w", updateErr) - } - return nil + if !errors.Is(err, ErrHandleAlreadyTaken) { + return err + } + // User exists but handle changed - update it + user, err = s.UpdateHandle(ctx, did, handle) + if err != nil { + return fmt.Errorf("failed to update handle for existing user: %w", err) } - return err } + // Best-effort and asynchronous: never fails or blocks indexing (a dead or + // slow PDS must not stall a login callback or firehose event, and must not + // dead-letter the post/comment event that triggered this index). + s.maybeBackfillProfile(ctx, user) + return nil } +// maybeBackfillProfile reconciles a missing profile for a freshly indexed user. A +// profile record fires exactly one firehose event (rkey "self" is written once); if +// that event was missed — e.g. relay throttling — nothing ever replays it, so the +// AppView must reconcile by reading the record directly from the user's own PDS. +// Only users with a completely empty profile are touched: for anyone else the +// firehose is the source of truth and a fetch could race a newer update. +// +// The emptiness check runs synchronously (the common no-op paths spawn nothing); +// when a fetch is actually needed, the fetch+store runs in a detached goroutine — +// decoupled from the caller's context via context.WithoutCancel with its own +// profileBackfillTimeout deadline — so a slow or dead PDS never blocks the +// IndexUser caller (OAuth login callback, Jetstream consumers) and the fetch is +// not killed when the triggering request context ends. +func (s *userService) maybeBackfillProfile(ctx context.Context, user *User) { + if s.profileBackfillClient == nil || user == nil { + return + } + if user.DisplayName != "" || user.Bio != "" || user.AvatarCID != "" || user.BannerCID != "" { + return // profile already indexed — firehose keeps it current + } + + go s.backfillProfile(context.WithoutCancel(ctx), user.DID, user.PDSURL) +} + +// backfillProfile is the detached best-effort fetch+store behind +// maybeBackfillProfile. Failures are logged only — for firehose-discovered users +// there is no automatic retry; run cmd/backfill-profiles to reconcile. +func (s *userService) backfillProfile(ctx context.Context, did, pdsURL string) { + ctx, cancel := context.WithTimeout(ctx, profileBackfillTimeout) + defer cancel() + + input, err := FetchProfileRecord(ctx, s.profileBackfillClient, pdsURL, did) + if err != nil { + slog.Warn("profile backfill: failed to fetch profile record (run cmd/backfill-profiles to reconcile)", + slog.String("did", did), + slog.String("pds_url", pdsURL), + slog.String("error", err.Error())) + return + } + if input == nil { + return // user has no profile record — nothing to apply + } + + // Re-check emptiness before writing: a concurrent firehose profile event may + // have landed while we were fetching, and the firehose is the source of truth. + current, err := s.userRepo.GetByDID(ctx, did) + if err != nil { + slog.Warn("profile backfill: failed to re-check user before storing fetched profile", + slog.String("did", did), + slog.String("error", err.Error())) + return + } + if current.DisplayName != "" || current.Bio != "" || current.AvatarCID != "" || current.BannerCID != "" { + return // firehose delivered a profile while we were fetching — keep it + } + + if _, err := s.userRepo.UpdateProfile(ctx, did, *input); err != nil { + slog.Warn("profile backfill: failed to store fetched profile", + slog.String("did", did), + slog.String("error", err.Error())) + return + } + + slog.Info("profile backfill: hydrated profile from PDS", + slog.String("did", did), + slog.String("pds_url", pdsURL)) +} + // GetProfile retrieves a user's full profile with aggregated statistics. // Returns a ProfileViewDetailed matching the social.coves.actor.defs#profileViewDetailed lexicon. // Avatar and Banner CIDs are transformed to URLs using the user's PDS URL. diff --git a/internal/db/postgres/discover_repo.go b/internal/db/postgres/discover_repo.go index ab74a5c..bb47754 100644 --- a/internal/db/postgres/discover_repo.go +++ b/internal/db/postgres/discover_repo.go @@ -36,30 +36,12 @@ func (r *postgresDiscoverRepo) GetDiscover(ctx context.Context, req discover.Get return nil, nil, discover.ErrInvalidCursor } - // Build the main query + // Build the main query (shared column list — see feedPostSelectClause) var selectClause string if req.Sort == "hot" { - selectClause = fmt.Sprintf(` - SELECT - 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.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`, feedHotRankExpression) + selectClause = feedPostSelectClause(feedHotRankExpression) } else { - selectClause = ` - SELECT - 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.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` + selectClause = feedPostSelectClause("NULL::numeric") } // Build optional viewer block filter (only when authenticated viewer is present) diff --git a/internal/db/postgres/feed_repo.go b/internal/db/postgres/feed_repo.go index b976d48..1c4f0f9 100644 --- a/internal/db/postgres/feed_repo.go +++ b/internal/db/postgres/feed_repo.go @@ -37,31 +37,13 @@ func (r *postgresFeedRepo) GetCommunityFeed(ctx context.Context, req communityFe return nil, nil, communityFeeds.ErrInvalidCursor } - // Build the main query + // Build the main query (shared column list — see feedPostSelectClause) // For hot sort, we need to compute and return the hot_rank for cursor building var selectClause string if req.Sort == "hot" { - selectClause = fmt.Sprintf(` - SELECT - 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.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`, feedHotRankExpression) + selectClause = feedPostSelectClause(feedHotRankExpression) } else { - selectClause = ` - SELECT - 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.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` + selectClause = feedPostSelectClause("NULL::numeric") } // Build optional viewer block filter (only when authenticated viewer is present) diff --git a/internal/db/postgres/feed_repo_base.go b/internal/db/postgres/feed_repo_base.go index df4529f..004dde6 100644 --- a/internal/db/postgres/feed_repo_base.go +++ b/internal/db/postgres/feed_repo_base.go @@ -6,14 +6,10 @@ import ( "database/sql" "encoding/base64" "encoding/hex" - "encoding/json" "fmt" - "log/slog" "strings" "time" - "Coves/internal/core/blobs" - "Coves/internal/core/communities" "Coves/internal/core/posts" ) @@ -354,129 +350,32 @@ func (r *feedRepoBase) buildCursor(post *posts.PostView, sort string, hotRank fl return base64.StdEncoding.EncodeToString([]byte(signed)) } -// scanFeedPost scans a database row into a PostView -// This is the shared scanning logic used by both timeline and discover feeds +// feedPostSelectClause builds the SELECT ... FROM posts p clause shared by the +// community, timeline, and discover feed queries. It is postViewSelectColumns (the +// single source of truth for PostView hydration, see post_repo.go) plus a hot_rank +// column: pass feedHotRankExpression for hot sort, or "NULL::numeric" for other sorts. +func feedPostSelectClause(hotRankExpr string) string { + return ` + SELECT` + postViewSelectColumns + `, + ` + hotRankExpr + ` as hot_rank + FROM posts p` +} + +// scanFeedPost scans a database row into a PostView plus its computed hot_rank. +// Rows MUST be selected via feedPostSelectClause; all PostView hydration is delegated +// to scanPostView so feed and post queries can never drift apart. func (r *feedRepoBase) scanFeedPost(rows *sql.Rows) (*posts.PostView, float64, error) { - var ( - postView posts.PostView - authorView posts.AuthorView - communityRef posts.CommunityRef - title, content sql.NullString - facets, embed sql.NullString - labelsJSON sql.NullString - editedAt sql.NullTime - communityHandle sql.NullString - communityAvatar sql.NullString - communityPDSURL sql.NullString - hotRank sql.NullFloat64 - ) - - err := rows.Scan( - &postView.URI, &postView.CID, &postView.RKey, - &authorView.DID, &authorView.Handle, - &communityRef.DID, &communityHandle, &communityRef.Name, &communityAvatar, &communityPDSURL, - &title, &content, &facets, &embed, &labelsJSON, - &postView.CreatedAt, &editedAt, &postView.IndexedAt, - &postView.UpvoteCount, &postView.DownvoteCount, &postView.Score, &postView.CommentCount, - &hotRank, - ) + var hotRank sql.NullFloat64 + postView, err := scanPostView(rows, &hotRank) if err != nil { return nil, 0, err } - // Build author view - postView.Author = &authorView - - // Build community ref - if communityHandle.Valid { - communityRef.Handle = communityHandle.String - } - // Hydrate avatar CID to URL using image proxy config (avatar_small preset for feed lists) - if avatarURL := blobs.HydrateImageURL(communities.GetImageProxyConfig(), communityPDSURL.String, communityRef.DID, communityAvatar.String, "avatar_small"); avatarURL != "" { - communityRef.Avatar = &avatarURL - } - if communityPDSURL.Valid { - communityRef.PDSURL = communityPDSURL.String - } - postView.Community = &communityRef - - // Parse facets JSON into local variable (will be added to record below) - // Log errors but continue - a single malformed post shouldn't break the entire feed - var facetArray []interface{} - if facets.Valid { - if err := json.Unmarshal([]byte(facets.String), &facetArray); err != nil { - slog.Warn("[FEED] failed to parse facets JSON", - "post_uri", postView.URI, - "error", err, - ) - } - } - - // Parse embed JSON - // Log errors but continue - a single malformed post shouldn't break the entire feed - if embed.Valid { - var embedData interface{} - if err := json.Unmarshal([]byte(embed.String), &embedData); err != nil { - slog.Warn("[FEED] failed to parse embed JSON", - "post_uri", postView.URI, - "error", err, - ) - } else { - postView.Embed = embedData - } - } - - // Build stats - postView.Stats = &posts.PostStats{ - Upvotes: postView.UpvoteCount, - Downvotes: postView.DownvoteCount, - Score: postView.Score, - CommentCount: postView.CommentCount, - } - - // Build the record (required by lexicon) - record := map[string]interface{}{ - "$type": "social.coves.community.post", - "community": communityRef.DID, - "author": authorView.DID, - "createdAt": postView.CreatedAt.Format(time.RFC3339), - } - - // Add optional fields to record if present - if title.Valid { - record["title"] = title.String - } - if content.Valid { - record["content"] = content.String - } - // Add facets to record if present - if facetArray != nil { - record["facets"] = facetArray - } - if postView.Embed != nil { - record["embed"] = postView.Embed - } - if labelsJSON.Valid { - // Labels are stored as JSONB containing full com.atproto.label.defs#selfLabels structure - // Deserialize and include in record - var selfLabels posts.SelfLabels - if err := json.Unmarshal([]byte(labelsJSON.String), &selfLabels); err != nil { - slog.Warn("[FEED] failed to parse labels JSON", - "post_uri", postView.URI, - "error", err, - ) - } else { - record["labels"] = selfLabels - } - } - - postView.Record = record - - // Return the computed hot_rank (0.0 if NULL for non-hot sorts) + // hot_rank is NULL for non-hot sorts hotRankValue := 0.0 if hotRank.Valid { hotRankValue = hotRank.Float64 } - return &postView, hotRankValue, nil + return postView, hotRankValue, nil } diff --git a/internal/db/postgres/post_repo.go b/internal/db/postgres/post_repo.go index f6adb3a..ef411e1 100644 --- a/internal/db/postgres/post_repo.go +++ b/internal/db/postgres/post_repo.go @@ -23,8 +23,9 @@ type postgresPostRepo struct { // postViewSelectColumns is the ordered SELECT list for hydrating a posts.PostView row. // 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. +// source of truth shared by GetViewsByURIs, GetByAuthor, and (via feedPostSelectClause) +// the community/timeline/discover feed queries, so the column order can only be defined +// in one place and the 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 @@ -32,7 +33,7 @@ type postgresPostRepo struct { // it is selected as-is. const postViewSelectColumns = ` p.uri, p.cid, p.rkey, - p.author_did, u.handle as author_handle, + p.author_did, u.handle as author_handle, u.display_name as author_display_name, u.avatar_cid as author_avatar, u.pds_url as author_pds_url, 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, @@ -158,7 +159,8 @@ func (r *postgresPostRepo) GetByURI(ctx context.Context, uri string) (*posts.Pos // GetViewsByURIs retrieves full post views for a set of canonical (DID-based) AT-URIs. // Returns a map keyed by URI; URIs that are missing or soft-deleted are simply absent // from the map (the caller emits notFoundPost markers for those). -// Reuses scanPostView for row scanning, so the SELECT column order must match GetByAuthor. +// Row scanning goes through the shared scanPostView, whose Scan order is kept aligned +// with the single source of truth for the SELECT list, postViewSelectColumns. // Backs the social.coves.community.post.get endpoint (feed hydration + permalinks). func (r *postgresPostRepo) GetViewsByURIs(ctx context.Context, uris []string) (map[string]*posts.PostView, error) { result := make(map[string]*posts.PostView, len(uris)) @@ -188,7 +190,7 @@ func (r *postgresPostRepo) GetViewsByURIs(ctx context.Context, uris []string) (m }() for rows.Next() { - postView, err := r.scanPostView(rows) + postView, err := scanPostView(rows) if err != nil { return nil, fmt.Errorf("failed to scan post: %w", err) } @@ -284,7 +286,7 @@ func (r *postgresPostRepo) GetByAuthor(ctx context.Context, req posts.GetAuthorP // Scan results var postViews []*posts.PostView for rows.Next() { - postView, err := r.scanPostView(rows) + postView, err := scanPostView(rows) if err != nil { return nil, nil, fmt.Errorf("failed to scan author post: %w", err) } @@ -378,36 +380,49 @@ func (r *postgresPostRepo) SoftDelete(ctx context.Context, uri string) error { return nil } -// scanPostView scans a database row into a PostView. Shared by GetByAuthor and -// GetViewsByURIs, which both SELECT postViewSelectColumns; the Scan order below MUST -// stay byte-aligned with that column list. -func (r *postgresPostRepo) scanPostView(rows *sql.Rows) (*posts.PostView, error) { +// scanPostView scans a database row into a PostView. Shared by GetByAuthor, +// GetViewsByURIs, and the feed repositories (via scanFeedPost), which all SELECT +// postViewSelectColumns; the Scan order below MUST stay byte-aligned with that column +// list. extraDest receives any columns a caller appends after postViewSelectColumns +// (the feed queries append hot_rank). +func scanPostView(rows *sql.Rows, extraDest ...interface{}) (*posts.PostView, error) { var ( - postView posts.PostView - authorView posts.AuthorView - communityRef posts.CommunityRef - title, content sql.NullString - facets, embed sql.NullString - labelsJSON sql.NullString - editedAt sql.NullTime - communityHandle sql.NullString - communityAvatar sql.NullString - communityPDSURL sql.NullString + postView posts.PostView + authorView posts.AuthorView + communityRef posts.CommunityRef + title, content sql.NullString + facets, embed sql.NullString + labelsJSON sql.NullString + editedAt sql.NullTime + authorDisplayName sql.NullString + authorAvatar sql.NullString + authorPDSURL sql.NullString + communityHandle sql.NullString + communityAvatar sql.NullString + communityPDSURL sql.NullString ) - err := rows.Scan( + dest := []interface{}{ &postView.URI, &postView.CID, &postView.RKey, - &authorView.DID, &authorView.Handle, + &authorView.DID, &authorView.Handle, &authorDisplayName, &authorAvatar, &authorPDSURL, &communityRef.DID, &communityHandle, &communityRef.Name, &communityAvatar, &communityPDSURL, &title, &content, &facets, &embed, &labelsJSON, &postView.CreatedAt, &editedAt, &postView.IndexedAt, &postView.UpvoteCount, &postView.DownvoteCount, &postView.Score, &postView.CommentCount, - ) - if err != nil { + } + dest = append(dest, extraDest...) + + if err := rows.Scan(dest...); err != nil { return nil, err } - // Build author view + // Build author view with profile hydration (avatar_small preset, same as community avatars) + if authorDisplayName.Valid && authorDisplayName.String != "" { + authorView.DisplayName = &authorDisplayName.String + } + if avatarURL := blobs.HydrateImageURL(communities.GetImageProxyConfig(), authorPDSURL.String, authorView.DID, authorAvatar.String, "avatar_small"); avatarURL != "" { + authorView.Avatar = &avatarURL + } postView.Author = &authorView // Build community ref @@ -481,7 +496,7 @@ func (r *postgresPostRepo) scanPostView(rows *sql.Rows) (*posts.PostView, error) if facetArray != nil { record["facets"] = facetArray } - if embed.Valid { + if postView.Embed != nil { record["embed"] = postView.Embed } if labelsJSON.Valid { diff --git a/internal/db/postgres/timeline_repo.go b/internal/db/postgres/timeline_repo.go index d748aab..6cb0dd3 100644 --- a/internal/db/postgres/timeline_repo.go +++ b/internal/db/postgres/timeline_repo.go @@ -37,31 +37,13 @@ func (r *postgresTimelineRepo) GetTimeline(ctx context.Context, req timeline.Get return nil, nil, timeline.ErrInvalidCursor } - // Build the main query + // Build the main query (shared column list — see feedPostSelectClause) // For hot sort, we need to compute and return the hot_rank for cursor building var selectClause string if req.Sort == "hot" { - selectClause = fmt.Sprintf(` - SELECT - 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.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`, feedHotRankExpression) + selectClause = feedPostSelectClause(feedHotRankExpression) } else { - selectClause = ` - SELECT - 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.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` + selectClause = feedPostSelectClause("NULL::numeric") } // Join with community_subscriptions to get posts from subscribed communities diff --git a/tests/integration/author_avatar_hydration_test.go b/tests/integration/author_avatar_hydration_test.go new file mode 100644 index 0000000..664fa47 --- /dev/null +++ b/tests/integration/author_avatar_hydration_test.go @@ -0,0 +1,105 @@ +package integration + +import ( + "Coves/internal/core/communityFeeds" + "Coves/internal/core/posts" + "Coves/internal/db/postgres" + "context" + "fmt" + "net/url" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// TestAuthorProfileHydration verifies that post views hydrate the author's display +// name and avatar from the users table across every read path that shares +// postViewSelectColumns/scanPostView: batch get-by-URI, author feed, and the +// community feed (which the timeline and discover feeds share their scanner with). +// +// Regression test for the bug where feeds and post views only hydrated the +// community avatar and author cards were always bare even for fully indexed users. +func TestAuthorProfileHydration(t *testing.T) { + if testing.Short() { + t.Skip("Skipping integration test in short mode") + } + + db := setupTestDB(t) + t.Cleanup(func() { _ = db.Close() }) + + ctx := context.Background() + testID := uniqueTestID() + + communityDID, err := createFeedTestCommunity(db, ctx, fmt.Sprintf("avhydr-%s", testID), fmt.Sprintf("owner-%s.test", testID)) + require.NoError(t, err) + + authorDID := fmt.Sprintf("did:plc:avhydr%s", testID) + postURI := createTestPost(t, db, communityDID, authorDID, "Author hydration post", 1, time.Now().Add(-1*time.Hour)) + + // Give the author an indexed profile (as the profile firehose consumer would) + const avatarCID = "bafkreiauthoravatar" + const displayName = "Hydrated Author" + _, err = db.ExecContext(ctx, + `UPDATE users SET display_name = $1, avatar_cid = $2 WHERE did = $3`, + displayName, avatarCID, authorDID) + require.NoError(t, err) + + assertAuthorHydrated := func(t *testing.T, author *posts.AuthorView, path string) { + t.Helper() + require.NotNil(t, author, "%s: author view missing", path) + assert.Equal(t, authorDID, author.DID, path) + require.NotNil(t, author.DisplayName, "%s: author display name not hydrated", path) + assert.Equal(t, displayName, *author.DisplayName, path) + require.NotNil(t, author.Avatar, "%s: author avatar not hydrated", path) + assert.Contains(t, *author.Avatar, avatarCID, "%s: avatar URL must reference the avatar CID", path) + // Image proxy is disabled in the test env, so the URL is the direct PDS + // getBlob form with the DID query-escaped. Pinning the author's DID here + // catches regressions that hydrate the community's DID/PDS instead. + assert.Contains(t, *author.Avatar, url.QueryEscape(authorDID), "%s: avatar URL must reference the author's own DID", path) + } + + postRepo := postgres.NewPostRepository(db) + + t.Run("GetViewsByURIs", func(t *testing.T) { + views, err := postRepo.GetViewsByURIs(ctx, []string{postURI}) + require.NoError(t, err) + require.Contains(t, views, postURI) + assertAuthorHydrated(t, views[postURI].Author, "GetViewsByURIs") + }) + + t.Run("GetByAuthor", func(t *testing.T) { + views, _, err := postRepo.GetByAuthor(ctx, posts.GetAuthorPostsRequest{ActorDID: authorDID, Limit: 10}) + require.NoError(t, err) + require.Len(t, views, 1) + assertAuthorHydrated(t, views[0].Author, "GetByAuthor") + }) + + t.Run("CommunityFeed", func(t *testing.T) { + feedRepo := postgres.NewCommunityFeedRepository(db, "test-cursor-secret") + for _, sort := range []string{"new", "hot"} { + feed, _, err := feedRepo.GetCommunityFeed(ctx, communityFeeds.GetCommunityFeedRequest{ + Community: communityDID, + Sort: sort, + Limit: 10, + }) + require.NoError(t, err, "sort=%s", sort) + require.Len(t, feed, 1, "sort=%s", sort) + assertAuthorHydrated(t, feed[0].Post.Author, "CommunityFeed sort="+sort) + } + }) + + t.Run("AuthorWithoutProfileStaysBare", func(t *testing.T) { + bareDID := fmt.Sprintf("did:plc:bare%s", testID) + bareURI := createTestPost(t, db, communityDID, bareDID, "Bare author post", 1, time.Now()) + + views, err := postRepo.GetViewsByURIs(ctx, []string{bareURI}) + require.NoError(t, err) + require.Contains(t, views, bareURI) + author := views[bareURI].Author + require.NotNil(t, author) + assert.Nil(t, author.DisplayName, "user without profile must not get a fabricated display name") + assert.Nil(t, author.Avatar, "user without profile must not get a fabricated avatar URL") + }) +} -- 2.51.2