diff --git a/appview/database/migrations/0002_comments_reply_root_rkey.sql b/appview/database/migrations/0002_comments_reply_root_rkey.sql new file mode 100644 index 0000000..b9e407e --- /dev/null +++ b/appview/database/migrations/0002_comments_reply_root_rkey.sql @@ -0,0 +1,13 @@ +-- migrate:no-transaction +-- +-- Supports the paginated GetCommentThread query: +-- WHERE reply_root = $1 AND rkey > $2 ORDER BY rkey ASC LIMIT $3 +-- The pre-existing idx_comments_reply_root (reply_root) would force a sort +-- pass over every matching row once thread sizes grow; the composite index +-- lets Postgres both filter and order using the index. +-- +-- CONCURRENTLY avoids blocking writes on the comments table. IF NOT EXISTS +-- makes the migration idempotent if a DBA created the index by hand already. + +CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_comments_reply_root_rkey + ON comments (reply_root, rkey); diff --git a/appview/database/migrations_test.go b/appview/database/migrations_test.go index dc1648b..df28814 100644 --- a/appview/database/migrations_test.go +++ b/appview/database/migrations_test.go @@ -146,4 +146,32 @@ func TestEmbeddedBaselineParses(t *testing.T) { if ms[0].Checksum == "" { t.Fatal("baseline checksum is empty") } + if ms[0].NoTransaction { + t.Fatal("baseline should run inside a transaction") + } +} + +// TestCommentsReplyRootRkeyMigrationIsConcurrent locks in the contract that +// the composite index migration opts out of the transaction wrapper. +// Running CREATE INDEX CONCURRENTLY inside a transaction is a Postgres error, +// so regressing this directive silently breaks a deploy. +func TestCommentsReplyRootRkeyMigrationIsConcurrent(t *testing.T) { + t.Parallel() + ms, err := loadMigrations(embeddedMigrations) + if err != nil { + t.Fatalf("embedded migrations did not load: %v", err) + } + var found *migration + for i := range ms { + if ms[i].Version == "0002" { + found = &ms[i] + break + } + } + if found == nil { + t.Fatal("expected migration 0002 to exist") + } + if !found.NoTransaction { + t.Fatalf("migration %s must carry the no-transaction directive (CREATE INDEX CONCURRENTLY)", found.Name) + } } diff --git a/appview/handlers/comment.go b/appview/handlers/comment.go index ed18e18..39119e5 100644 --- a/appview/handlers/comment.go +++ b/appview/handlers/comment.go @@ -106,8 +106,9 @@ func (h *Handlers) GetCommentThread(c echo.Context) error { } reqDID := requestingDID(c) + ctx := c.Request().Context() - rootQuery := h.db.WithContext(c.Request().Context()).Where("at_uri = ?", uri) + rootQuery := h.db.WithContext(ctx).Where("at_uri = ?", uri) rootQuery = h.excludeBlockedDIDs(rootQuery, reqDID, "did") var root database.Comment @@ -115,7 +116,20 @@ func (h *Handlers) GetCommentThread(c echo.Context) error { return writeError(c, http.StatusNotFound, "NotFound", "comment not found") } - repliesQuery := h.db.WithContext(c.Request().Context()).Where("reply_root = ?", root.ATURI).Order("created_at ASC") + limit := parseLimit(c.QueryParam("limit"), 50, 200) + cursor := c.QueryParam("cursor") + + // Order and cursor both by rkey. AT Protocol TIDs are monotonically + // time-sortable, so this gives chronological order without the skip / + // duplicate hazard that a cursor on rkey with ORDER BY created_at would + // create when two replies share a timestamp. + repliesQuery := h.db.WithContext(ctx). + Where("reply_root = ?", root.ATURI). + Order("rkey ASC"). + Limit(limit + 1) + if cursor != "" { + repliesQuery = repliesQuery.Where("rkey > ?", cursor) + } repliesQuery = h.excludeBlockedDIDs(repliesQuery, reqDID, "did") var replies []database.Comment @@ -123,6 +137,23 @@ func (h *Handlers) GetCommentThread(c echo.Context) error { return h.internalError(c, "GetCommentThread.replies", err) } + nextCursor := "" + if len(replies) > limit { + nextCursor = replies[limit-1].Rkey + replies = replies[:limit] + } + + // reply_count is the raw thread size — it does not account for the + // caller's block filter, so the number can exceed the replies they + // actually see. Clients should treat it as an upper bound. + var replyCount int64 + if err := h.db.WithContext(ctx). + Model(&database.Comment{}). + Where("reply_root = ?", root.ATURI). + Count(&replyCount).Error; err != nil { + return h.internalError(c, "GetCommentThread.count", err) + } + replyNodes := make([]map[string]any, 0, len(replies)) for _, row := range replies { replyNodes = append(replyNodes, map[string]any{ @@ -136,5 +167,7 @@ func (h *Handlers) GetCommentThread(c echo.Context) error { "comment": toCommentResponse(root), "replies": replyNodes, }, + "reply_count": replyCount, + "cursor": nextCursor, }) } diff --git a/appview/handlers/inbox.go b/appview/handlers/inbox.go index 6f8bf82..e01e938 100644 --- a/appview/handlers/inbox.go +++ b/appview/handlers/inbox.go @@ -4,10 +4,19 @@ import ( "net/http" "strconv" - "tangled.org/sparrowtek.com/effem-AppView/appview/database" "github.com/labstack/echo/v4" + "tangled.org/sparrowtek.com/effem-AppView/appview/database" + "tangled.org/sparrowtek.com/effem-AppView/appview/metrics" ) +// inboxSubscriptionCap limits how many subscriptions GetInbox considers on a +// single request. Every subscription turns into one Podcast Index fetch (or +// cache read) so the work per request grows linearly; without a cap a user +// with hundreds of subs could turn a single inbox call into hundreds of +// upstream calls. Most-recently subscribed feeds are used first on the +// assumption that those are the ones the user actually listens to. +const inboxSubscriptionCap = 200 + func (h *Handlers) GetInbox(c echo.Context) error { did := c.QueryParam("did") if did == "" { @@ -24,9 +33,14 @@ func (h *Handlers) GetInbox(c echo.Context) error { } var subs []database.Subscription - if err := h.db.WithContext(c.Request().Context()).Where("did = ?", did).Find(&subs).Error; err != nil { + if err := h.db.WithContext(c.Request().Context()). + Where("did = ?", did). + Order("rkey DESC"). + Limit(inboxSubscriptionCap). + Find(&subs).Error; err != nil { return h.internalError(c, "GetInbox.subscriptions", err) } + metrics.ObserveInboxSubscriptions(len(subs)) if len(subs) == 0 { return c.JSON(http.StatusOK, map[string]any{"items": []any{}, "cursor": ""}) } diff --git a/appview/metrics/metrics.go b/appview/metrics/metrics.go index 37d425e..f526388 100644 --- a/appview/metrics/metrics.go +++ b/appview/metrics/metrics.go @@ -98,6 +98,17 @@ var ( }, []string{"kind"}, ) + + inboxSubscriptionCount = prometheus.NewHistogram( + prometheus.HistogramOpts{ + Namespace: namespace, + Name: "inbox_subscription_count", + Help: "Number of subscriptions considered per GetInbox request (post-cap).", + // Powers of 2 from 1 to 512: 1,2,4,8,16,32,64,128,256,512. + // The cap lives at 200 today; buckets above that catch future cap bumps. + Buckets: prometheus.ExponentialBuckets(1, 2, 10), + }, + ) ) func init() { @@ -110,6 +121,7 @@ func init() { piCacheHitsTotal, piCacheMissesTotal, rateLimitRejectionsTotal, + inboxSubscriptionCount, ) } @@ -209,3 +221,10 @@ func IncCacheMiss(endpoint string) { func IncRateLimitRejection(kind string) { rateLimitRejectionsTotal.WithLabelValues(kind).Inc() } + +// ObserveInboxSubscriptions records how many subscriptions were considered +// on a single GetInbox request. Use this to decide whether the current cap +// is meeting real-world usage. +func ObserveInboxSubscriptions(n int) { + inboxSubscriptionCount.Observe(float64(n)) +}