diff --git a/appview/database/migrations.go b/appview/database/migrations.go index 6a14421..d033510 100644 --- a/appview/database/migrations.go +++ b/appview/database/migrations.go @@ -4,12 +4,15 @@ import "gorm.io/gorm" func RunMigrations(db *gorm.DB) error { return db.AutoMigrate( - &FirehoseCursor{}, &Subscription{}, &Comment{}, &Recommendation{}, &PodcastList{}, &Bookmark{}, &Profile{}, + &FirehoseCursor{}, + &PICache{}, + &PodcastStats{}, + &EpisodeStats{}, ) } diff --git a/appview/database/models.go b/appview/database/models.go index 3c12d1e..0ad60d3 100644 --- a/appview/database/models.go +++ b/appview/database/models.go @@ -2,28 +2,30 @@ package database import "time" -const DefaultCursorName = "main" - type FirehoseCursor struct { - Name string `gorm:"primaryKey;size:64"` + ID uint `gorm:"primaryKey"` Seq int64 `gorm:"not null"` UpdatedAt time.Time `gorm:"autoUpdateTime"` } +func (FirehoseCursor) TableName() string { return "firehose_cursor" } + type Subscription struct { ID uint `gorm:"primaryKey"` - DID string `gorm:"size:255;not null;index:idx_subscriptions_did_rkey,unique"` + DID string `gorm:"size:255;not null;index:idx_subscriptions_did_rkey,unique;index:idx_subscriptions_did"` Rkey string `gorm:"size:512;not null;index:idx_subscriptions_did_rkey,unique"` - FeedID int `gorm:"not null;index"` + FeedID int `gorm:"not null;index:idx_subscriptions_feed_id"` FeedURL string `gorm:"size:2048"` PodcastGuid string `gorm:"size:512"` CreatedAt string `gorm:"size:64;not null"` IndexedAt time.Time `gorm:"autoCreateTime"` } +func (Subscription) TableName() string { return "subscriptions" } + type Comment struct { ID uint `gorm:"primaryKey"` - DID string `gorm:"size:255;not null;index:idx_comments_did_rkey,unique"` + DID string `gorm:"size:255;not null;index:idx_comments_did_rkey,unique;index:idx_comments_did"` Rkey string `gorm:"size:512;not null;index:idx_comments_did_rkey,unique"` ATURI string `gorm:"size:1024;index;not null"` FeedID int `gorm:"not null;index:idx_comments_episode"` @@ -32,17 +34,18 @@ type Comment struct { PodcastGuid string `gorm:"size:512"` Text string `gorm:"type:text;not null"` TimestampS *int `gorm:"index"` - RootURI string `gorm:"size:1024;index"` - RootCID string `gorm:"size:255"` - ParentURI string `gorm:"size:1024;index"` - ParentCID string `gorm:"size:255"` + ReplyRoot string `gorm:"size:1024;index:idx_comments_reply_root"` + ReplyParent string `gorm:"size:1024"` + Facets []byte `gorm:"type:jsonb"` CreatedAt string `gorm:"size:64;not null"` IndexedAt time.Time `gorm:"autoCreateTime"` } +func (Comment) TableName() string { return "comments" } + type Recommendation struct { ID uint `gorm:"primaryKey"` - DID string `gorm:"size:255;not null;index:idx_recommendations_did_rkey,unique"` + DID string `gorm:"size:255;not null;index:idx_recommendations_did_rkey,unique;index:idx_recommendations_did"` Rkey string `gorm:"size:512;not null;index:idx_recommendations_did_rkey,unique"` FeedID int `gorm:"not null;index:idx_recommendations_episode"` EpisodeID int `gorm:"not null;index:idx_recommendations_episode"` @@ -53,9 +56,11 @@ type Recommendation struct { IndexedAt time.Time `gorm:"autoCreateTime"` } +func (Recommendation) TableName() string { return "recommendations" } + type PodcastList struct { ID uint `gorm:"primaryKey"` - DID string `gorm:"size:255;not null;index:idx_lists_did_rkey,unique"` + DID string `gorm:"size:255;not null;index:idx_lists_did_rkey,unique;index:idx_lists_did"` Rkey string `gorm:"size:512;not null;index:idx_lists_did_rkey,unique"` Name string `gorm:"size:500;not null"` Description string `gorm:"type:text"` @@ -64,9 +69,11 @@ type PodcastList struct { IndexedAt time.Time `gorm:"autoCreateTime"` } +func (PodcastList) TableName() string { return "podcast_lists" } + type Bookmark struct { ID uint `gorm:"primaryKey"` - DID string `gorm:"size:255;not null;index:idx_bookmarks_did_rkey,unique"` + DID string `gorm:"size:255;not null;index:idx_bookmarks_did_rkey,unique;index:idx_bookmarks_did"` Rkey string `gorm:"size:512;not null;index:idx_bookmarks_did_rkey,unique"` FeedID int `gorm:"not null;index:idx_bookmarks_episode"` EpisodeID int `gorm:"not null;index:idx_bookmarks_episode"` @@ -77,11 +84,46 @@ type Bookmark struct { IndexedAt time.Time `gorm:"autoCreateTime"` } +func (Bookmark) TableName() string { return "bookmarks" } + type Profile struct { - ID uint `gorm:"primaryKey"` - DID string `gorm:"size:255;not null;uniqueIndex"` + DID string `gorm:"primaryKey;size:255"` DisplayName string `gorm:"size:640"` Description string `gorm:"type:text"` FavoriteGenres []byte `gorm:"type:jsonb"` IndexedAt time.Time `gorm:"autoCreateTime"` } + +func (Profile) TableName() string { return "profiles" } + +type PICache struct { + ID uint `gorm:"primaryKey"` + CacheKey string `gorm:"uniqueIndex;size:512"` + Response []byte `gorm:"type:jsonb;not null"` + ExpiresAt time.Time `gorm:"index:idx_pi_cache_expires;not null"` + CreatedAt time.Time `gorm:"autoCreateTime"` + UpdatedAt time.Time `gorm:"autoUpdateTime"` +} + +func (PICache) TableName() string { return "pi_cache" } + +type PodcastStats struct { + FeedID int `gorm:"primaryKey"` + SubscriberCount int `gorm:"default:0"` + CommentCount int `gorm:"default:0"` + RecommendationCount int `gorm:"default:0"` + LastUpdated time.Time `gorm:"autoUpdateTime"` +} + +func (PodcastStats) TableName() string { return "podcast_stats" } + +type EpisodeStats struct { + EpisodeID int `gorm:"primaryKey"` + FeedID int `gorm:"not null;index:idx_episode_stats_feed"` + CommentCount int `gorm:"default:0"` + RecommendationCount int `gorm:"default:0"` + BookmarkCount int `gorm:"default:0"` + LastUpdated time.Time `gorm:"autoUpdateTime"` +} + +func (EpisodeStats) TableName() string { return "episode_stats" } diff --git a/appview/firehose.go b/appview/firehose.go index 58cd814..8fb6e5e 100644 --- a/appview/firehose.go +++ b/appview/firehose.go @@ -19,7 +19,6 @@ import ( "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" "github.com/gorilla/websocket" - "gorm.io/gorm/clause" ) const effemNSPrefix = "xyz.effem." @@ -163,7 +162,7 @@ func (srv *Server) persistCursorLoop(ctx context.Context) { func (srv *Server) loadFirehoseCursor() (int64, error) { var cur database.FirehoseCursor - err := srv.db.Where("name = ?", database.DefaultCursorName).First(&cur).Error + err := srv.db.First(&cur).Error if err != nil { return 0, err } @@ -171,12 +170,12 @@ func (srv *Server) loadFirehoseCursor() (int64, error) { } func (srv *Server) saveFirehoseCursor(seq int64) error { - cur := database.FirehoseCursor{ - Name: database.DefaultCursorName, - Seq: seq, + var cur database.FirehoseCursor + err := srv.db.First(&cur).Error + if err != nil { + cur = database.FirehoseCursor{Seq: seq} + return srv.db.Create(&cur).Error } - return srv.db.Clauses(clause.OnConflict{ - Columns: []clause.Column{{Name: "name"}}, - DoUpdates: clause.AssignmentColumns([]string{"seq", "updated_at"}), - }).Create(&cur).Error + cur.Seq = seq + return srv.db.Save(&cur).Error } diff --git a/appview/handlers/comment.go b/appview/handlers/comment.go index e3a798f..1c5cd47 100644 --- a/appview/handlers/comment.go +++ b/appview/handlers/comment.go @@ -54,7 +54,7 @@ func (h *Handlers) GetCommentThread(c echo.Context) error { } var replies []database.Comment - if err := h.db.WithContext(c.Request().Context()).Where("root_uri = ?", root.ATURI).Order("created_at ASC").Find(&replies).Error; err != nil { + if err := h.db.WithContext(c.Request().Context()).Where("reply_root = ?", root.ATURI).Order("created_at ASC").Find(&replies).Error; err != nil { return writeError(c, http.StatusInternalServerError, "InternalError", "") } diff --git a/appview/handlers/episodes.go b/appview/handlers/episodes.go index 776e864..72b8877 100644 --- a/appview/handlers/episodes.go +++ b/appview/handlers/episodes.go @@ -13,12 +13,38 @@ func (h *Handlers) GetEpisodes(c echo.Context) error { } max := parseLimit(c.QueryParam("max"), 50, 200) - payload, err := h.pi.GetEpisodesByFeedID(feedID, max) + raw, err := h.pi.GetEpisodesByFeedID(feedID, max) if err != nil { h.logger.Warn("get episodes failed", "err", err, "feedId", feedID) return writeError(c, http.StatusBadGateway, "UpstreamError", "podcast index request failed") } - return writeJSONBlob(c, payload) + + payload, err := decodePayload(raw) + if err != nil { + return writeJSONBlob(c, raw) + } + + arr, ok := payload["items"].([]any) + if !ok { + return c.JSON(http.StatusOK, payload) + } + + items := make([]any, 0, len(arr)) + for _, v := range arr { + item, ok := v.(map[string]any) + if !ok { + items = append(items, v) + continue + } + episodeID := anyToInt(item["id"]) + if episodeID > 0 { + item["social"] = h.episodeSocialCounts(c.Request().Context(), feedID, episodeID) + } + items = append(items, item) + } + payload["items"] = items + + return c.JSON(http.StatusOK, payload) } func (h *Handlers) GetEpisode(c echo.Context) error { @@ -27,10 +53,19 @@ func (h *Handlers) GetEpisode(c echo.Context) error { return writeError(c, http.StatusBadRequest, "InvalidRequest", "episodeId is required") } - payload, err := h.pi.GetEpisodeByID(episodeID) + raw, err := h.pi.GetEpisodeByID(episodeID) if err != nil { h.logger.Warn("get episode failed", "err", err, "episodeId", episodeID) return writeError(c, http.StatusBadGateway, "UpstreamError", "podcast index request failed") } - return writeJSONBlob(c, payload) + + payload, err := decodePayload(raw) + if err != nil { + return writeJSONBlob(c, raw) + } + + feedID := parseInt(c.QueryParam("feedId"), 0) + payload["social"] = h.episodeSocialCounts(c.Request().Context(), feedID, episodeID) + + return c.JSON(http.StatusOK, payload) } diff --git a/appview/handlers/inbox.go b/appview/handlers/inbox.go index 812b407..b106adf 100644 --- a/appview/handlers/inbox.go +++ b/appview/handlers/inbox.go @@ -2,6 +2,7 @@ package handlers import ( "net/http" + "strconv" "github.com/SparrowTek/effem-appview/appview/database" "github.com/labstack/echo/v4" @@ -14,30 +15,87 @@ func (h *Handlers) GetInbox(c echo.Context) error { } limit := parseLimit(c.QueryParam("limit"), 50, 100) + offset := 0 + if cursor := c.QueryParam("cursor"); cursor != "" { + offset = parseInt(cursor, 0) + if offset < 0 { + offset = 0 + } + } var subs []database.Subscription if err := h.db.WithContext(c.Request().Context()).Where("did = ?", did).Find(&subs).Error; err != nil { return writeError(c, http.StatusInternalServerError, "InternalError", "") } - if len(subs) == 0 { - return c.JSON(http.StatusOK, map[string]any{"items": []any{}}) + return c.JSON(http.StatusOK, map[string]any{"items": []any{}, "cursor": ""}) } - feedIDs := make([]int, 0, len(subs)) - for _, s := range subs { - feedIDs = append(feedIDs, s.FeedID) + type seenKey struct { + FeedID int + EpisodeID int } + seen := map[seenKey]struct{}{} + items := make([]map[string]any, 0) - var items []database.Comment - err := h.db.WithContext(c.Request().Context()). - Where("feed_id IN ?", feedIDs). - Order("created_at DESC"). - Limit(limit). - Find(&items).Error - if err != nil { - return writeError(c, http.StatusInternalServerError, "InternalError", "") + for _, sub := range subs { + raw, err := h.pi.GetEpisodesByFeedID(sub.FeedID, 5) + if err != nil { + h.logger.Warn("inbox episodes fetch failed", "feedId", sub.FeedID, "err", err) + continue + } + payload, err := decodePayload(raw) + if err != nil { + continue + } + arr, ok := payload["items"].([]any) + if !ok { + continue + } + + for _, v := range arr { + item, ok := v.(map[string]any) + if !ok { + continue + } + episodeID := anyToInt(item["id"]) + if episodeID <= 0 { + continue + } + key := seenKey{FeedID: sub.FeedID, EpisodeID: episodeID} + if _, found := seen[key]; found { + continue + } + seen[key] = struct{}{} + + if feedIDFromItem(item) <= 0 { + item["feedId"] = sub.FeedID + } + + items = append(items, item) + } + } + + sortEpisodeItemsByPublishedDesc(items) + + if offset >= len(items) { + return c.JSON(http.StatusOK, map[string]any{"items": []any{}, "cursor": ""}) + } + + end := offset + limit + if end > len(items) { + end = len(items) + } + + nextCursor := "" + if end < len(items) { + nextCursor = strconv.Itoa(end) + } + + paged := make([]any, 0, end-offset) + for _, item := range items[offset:end] { + paged = append(paged, item) } - return c.JSON(http.StatusOK, map[string]any{"items": items}) + return c.JSON(http.StatusOK, map[string]any{"items": paged, "cursor": nextCursor}) } diff --git a/appview/handlers/list.go b/appview/handlers/list.go index 92af353..426e624 100644 --- a/appview/handlers/list.go +++ b/appview/handlers/list.go @@ -5,15 +5,26 @@ import ( "net/http" "github.com/SparrowTek/effem-appview/appview/database" + "github.com/bluesky-social/indigo/atproto/syntax" "github.com/labstack/echo/v4" ) func (h *Handlers) GetList(c echo.Context) error { - did := c.QueryParam("did") - rkey := c.QueryParam("rkey") - if did == "" || rkey == "" { - return writeError(c, http.StatusBadRequest, "InvalidRequest", "did and rkey are required") + uriRaw := c.QueryParam("uri") + if uriRaw == "" { + return writeError(c, http.StatusBadRequest, "InvalidRequest", "uri is required") + } + + uri, err := syntax.ParseATURI(uriRaw) + if err != nil { + return writeError(c, http.StatusBadRequest, "InvalidRequest", "invalid AT URI") } + if uri.Collection().String() != "xyz.effem.feed.list" || uri.RecordKey().String() == "" { + return writeError(c, http.StatusBadRequest, "InvalidRequest", "uri must reference xyz.effem.feed.list") + } + + did := uri.Authority().String() + rkey := uri.RecordKey().String() var list database.PodcastList if err := h.db.WithContext(c.Request().Context()).Where("did = ? AND rkey = ?", did, rkey).First(&list).Error; err != nil { diff --git a/appview/handlers/podcast.go b/appview/handlers/podcast.go index 56a90a0..e62c0a1 100644 --- a/appview/handlers/podcast.go +++ b/appview/handlers/podcast.go @@ -12,25 +12,51 @@ func (h *Handlers) GetPodcast(c echo.Context) error { return writeError(c, http.StatusBadRequest, "InvalidRequest", "feedId is required") } - payload, err := h.pi.GetPodcastByFeedID(feedID) + raw, err := h.pi.GetPodcastByFeedID(feedID) if err != nil { h.logger.Warn("get podcast failed", "err", err, "feedId", feedID) return writeError(c, http.StatusBadGateway, "UpstreamError", "podcast index request failed") } - return writeJSONBlob(c, payload) + + payload, err := decodePayload(raw) + if err != nil { + return writeJSONBlob(c, raw) + } + payload["social"] = h.podcastSocialCounts(c.Request().Context(), feedID) + + return c.JSON(http.StatusOK, payload) } func (h *Handlers) GetTrending(c echo.Context) error { max := parseLimit(c.QueryParam("max"), 20, 100) lang := c.QueryParam("lang") - categories := c.QueryParam("categories") + categories := c.QueryParam("cat") + if categories == "" { + categories = c.QueryParam("categories") + } - payload, err := h.pi.GetTrending(max, lang, categories) + raw, err := h.pi.GetTrending(max, lang, categories) if err != nil { h.logger.Warn("get trending failed", "err", err) return writeError(c, http.StatusBadGateway, "UpstreamError", "podcast index request failed") } - return writeJSONBlob(c, payload) + + payload, err := decodePayload(raw) + if err != nil { + return writeJSONBlob(c, raw) + } + + items, key := extractFeedItems(payload) + for _, item := range items { + feedID := feedIDFromItem(item) + if feedID <= 0 { + continue + } + item["social"] = h.podcastSocialCounts(c.Request().Context(), feedID) + } + storeFeedItems(payload, key, items) + + return c.JSON(http.StatusOK, payload) } func (h *Handlers) GetCategories(c echo.Context) error { diff --git a/appview/handlers/profile.go b/appview/handlers/profile.go index 986c4fa..60c22b5 100644 --- a/appview/handlers/profile.go +++ b/appview/handlers/profile.go @@ -24,10 +24,28 @@ func (h *Handlers) GetProfile(c echo.Context) error { _ = json.Unmarshal(profile.FavoriteGenres, &genres) } + var subscriptionCount int64 + _ = h.db.WithContext(c.Request().Context()).Model(&database.Subscription{}).Where("did = ?", did).Count(&subscriptionCount).Error + var commentCount int64 + _ = h.db.WithContext(c.Request().Context()).Model(&database.Comment{}).Where("did = ?", did).Count(&commentCount).Error + var recommendationCount int64 + _ = h.db.WithContext(c.Request().Context()).Model(&database.Recommendation{}).Where("did = ?", did).Count(&recommendationCount).Error + var bookmarkCount int64 + _ = h.db.WithContext(c.Request().Context()).Model(&database.Bookmark{}).Where("did = ?", did).Count(&bookmarkCount).Error + var listCount int64 + _ = h.db.WithContext(c.Request().Context()).Model(&database.PodcastList{}).Where("did = ?", did).Count(&listCount).Error + return c.JSON(http.StatusOK, map[string]any{ "did": profile.DID, "displayName": profile.DisplayName, "description": profile.Description, "favoriteGenres": genres, + "stats": map[string]any{ + "subscriptionCount": subscriptionCount, + "commentCount": commentCount, + "recommendationCount": recommendationCount, + "bookmarkCount": bookmarkCount, + "listCount": listCount, + }, }) } diff --git a/appview/handlers/recommendation.go b/appview/handlers/recommendation.go index 3c5d19c..ec1ac13 100644 --- a/appview/handlers/recommendation.go +++ b/appview/handlers/recommendation.go @@ -2,6 +2,7 @@ package handlers import ( "net/http" + "time" "github.com/SparrowTek/effem-appview/appview/database" "github.com/labstack/echo/v4" @@ -46,11 +47,27 @@ type popularRow struct { func (h *Handlers) GetPopular(c echo.Context) error { limit := parseLimit(c.QueryParam("limit"), 20, 100) - var rows []popularRow + period := c.QueryParam("period") + + q := h.db.WithContext(c.Request().Context()).Model(&database.Recommendation{}) + if period != "" { + now := time.Now().UTC() + var cutoff time.Time + switch period { + case "day": + cutoff = now.Add(-24 * time.Hour) + case "week": + cutoff = now.Add(-7 * 24 * time.Hour) + case "month": + cutoff = now.Add(-30 * 24 * time.Hour) + default: + return writeError(c, http.StatusBadRequest, "InvalidRequest", "period must be day, week, or month") + } + q = q.Where("created_at >= ?", cutoff.Format(time.RFC3339)) + } - err := h.db.WithContext(c.Request().Context()). - Model(&database.Recommendation{}). - Select("feed_id, episode_id, COUNT(*) as count"). + var rows []popularRow + err := q.Select("feed_id, episode_id, COUNT(*) as count"). Group("feed_id, episode_id"). Order("count DESC"). Limit(limit). diff --git a/appview/handlers/search.go b/appview/handlers/search.go index 71ee3e3..e7f4919 100644 --- a/appview/handlers/search.go +++ b/appview/handlers/search.go @@ -13,12 +13,28 @@ func (h *Handlers) SearchPodcasts(c echo.Context) error { } max := parseLimit(c.QueryParam("max"), 20, 100) - payload, err := h.pi.SearchByTerm(q, max) + raw, err := h.pi.SearchByTerm(q, max) if err != nil { h.logger.Warn("podcast search failed", "err", err) return writeError(c, http.StatusBadGateway, "UpstreamError", "podcast index request failed") } - return writeJSONBlob(c, payload) + + payload, err := decodePayload(raw) + if err != nil { + return writeJSONBlob(c, raw) + } + + items, key := extractFeedItems(payload) + for _, item := range items { + feedID := feedIDFromItem(item) + if feedID <= 0 { + continue + } + item["social"] = h.podcastSocialCounts(c.Request().Context(), feedID) + } + storeFeedItems(payload, key, items) + + return c.JSON(http.StatusOK, payload) } func (h *Handlers) SearchEpisodes(c echo.Context) error { diff --git a/appview/handlers/social.go b/appview/handlers/social.go new file mode 100644 index 0000000..046d3dc --- /dev/null +++ b/appview/handlers/social.go @@ -0,0 +1,178 @@ +package handlers + +import ( + "context" + "encoding/json" + "sort" + "strconv" + + "github.com/SparrowTek/effem-appview/appview/database" +) + +type SocialCounts struct { + SubscriberCount int `json:"subscriberCount"` + CommentCount int `json:"commentCount"` + RecommendationCount int `json:"recommendationCount"` + SubscribedByFollowing []string `json:"subscribedByFollowing"` +} + +type EpisodeSocialCounts struct { + CommentCount int `json:"commentCount"` + RecommendationCount int `json:"recommendationCount"` + BookmarkCount int `json:"bookmarkCount"` +} + +func anyToInt(v any) int { + switch x := v.(type) { + case int: + return x + case int8: + return int(x) + case int16: + return int(x) + case int32: + return int(x) + case int64: + return int(x) + case uint: + return int(x) + case uint8: + return int(x) + case uint16: + return int(x) + case uint32: + return int(x) + case uint64: + return int(x) + case float64: + return int(x) + case json.Number: + i, _ := x.Int64() + return int(i) + case string: + i, _ := strconv.Atoi(x) + return i + default: + return 0 + } +} + +func decodePayload(raw json.RawMessage) (map[string]any, error) { + out := map[string]any{} + err := json.Unmarshal(raw, &out) + return out, err +} + +func extractFeedItems(payload map[string]any) ([]map[string]any, string) { + keys := []string{"feeds", "items"} + for _, key := range keys { + arr, ok := payload[key].([]any) + if !ok { + continue + } + items := make([]map[string]any, 0, len(arr)) + for _, v := range arr { + m, ok := v.(map[string]any) + if ok { + items = append(items, m) + } + } + return items, key + } + return nil, "" +} + +func storeFeedItems(payload map[string]any, key string, items []map[string]any) { + if key == "" { + return + } + arr := make([]any, 0, len(items)) + for _, item := range items { + arr = append(arr, item) + } + payload[key] = arr +} + +func feedIDFromItem(item map[string]any) int { + if id := anyToInt(item["feedId"]); id > 0 { + return id + } + if id := anyToInt(item["id"]); id > 0 { + return id + } + return 0 +} + +func (h *Handlers) podcastSocialCounts(ctx context.Context, feedID int) SocialCounts { + if feedID <= 0 { + return SocialCounts{} + } + + var stats database.PodcastStats + if err := h.db.WithContext(ctx).Where("feed_id = ?", feedID).First(&stats).Error; err == nil { + return SocialCounts{ + SubscriberCount: stats.SubscriberCount, + CommentCount: stats.CommentCount, + RecommendationCount: stats.RecommendationCount, + SubscribedByFollowing: []string{}, + } + } + + var subCount int64 + _ = h.db.WithContext(ctx).Model(&database.Subscription{}).Where("feed_id = ?", feedID).Count(&subCount).Error + var commentCount int64 + _ = h.db.WithContext(ctx).Model(&database.Comment{}).Where("feed_id = ?", feedID).Count(&commentCount).Error + var recCount int64 + _ = h.db.WithContext(ctx).Model(&database.Recommendation{}).Where("feed_id = ?", feedID).Count(&recCount).Error + + return SocialCounts{ + SubscriberCount: int(subCount), + CommentCount: int(commentCount), + RecommendationCount: int(recCount), + SubscribedByFollowing: []string{}, + } +} + +func sortEpisodeItemsByPublishedDesc(items []map[string]any) { + sort.SliceStable(items, func(i, j int) bool { + left := anyToInt(items[i]["datePublished"]) + right := anyToInt(items[j]["datePublished"]) + if left == right { + return anyToInt(items[i]["id"]) > anyToInt(items[j]["id"]) + } + return left > right + }) +} + +func (h *Handlers) episodeSocialCounts(ctx context.Context, feedID, episodeID int) EpisodeSocialCounts { + if episodeID <= 0 { + return EpisodeSocialCounts{} + } + + var stats database.EpisodeStats + if err := h.db.WithContext(ctx).Where("episode_id = ?", episodeID).First(&stats).Error; err == nil { + return EpisodeSocialCounts{ + CommentCount: stats.CommentCount, + RecommendationCount: stats.RecommendationCount, + BookmarkCount: stats.BookmarkCount, + } + } + + q := h.db.WithContext(ctx) + if feedID > 0 { + q = q.Where("feed_id = ?", feedID) + } + + var commentCount int64 + _ = q.Model(&database.Comment{}).Where("episode_id = ?", episodeID).Count(&commentCount).Error + var recommendationCount int64 + _ = q.Model(&database.Recommendation{}).Where("episode_id = ?", episodeID).Count(&recommendationCount).Error + var bookmarkCount int64 + _ = q.Model(&database.Bookmark{}).Where("episode_id = ?", episodeID).Count(&bookmarkCount).Error + + return EpisodeSocialCounts{ + CommentCount: int(commentCount), + RecommendationCount: int(recommendationCount), + BookmarkCount: int(bookmarkCount), + } +} diff --git a/appview/handlers/subscription.go b/appview/handlers/subscription.go index 6de4889..cab4d9d 100644 --- a/appview/handlers/subscription.go +++ b/appview/handlers/subscription.go @@ -44,6 +44,39 @@ func (h *Handlers) GetSubscribers(c echo.Context) error { return writeError(c, http.StatusBadRequest, "InvalidRequest", "feedId is required") } + limit := parseLimit(c.QueryParam("limit"), 50, 100) + cursor := c.QueryParam("cursor") + + q := h.db.WithContext(c.Request().Context()). + Where("feed_id = ?", feedID). + Order("rkey DESC"). + Limit(limit + 1) + if cursor != "" { + q = q.Where("rkey < ?", cursor) + } + + var subs []database.Subscription + if err := q.Find(&subs).Error; err != nil { + return writeError(c, http.StatusInternalServerError, "InternalError", "") + } + + nextCursor := "" + if len(subs) > limit { + nextCursor = subs[limit-1].Rkey + subs = subs[:limit] + } + + type subscriber struct { + DID string `json:"did"` + Rkey string `json:"rkey"` + CreatedAt string `json:"createdAt"` + } + + list := make([]subscriber, 0, len(subs)) + for _, row := range subs { + list = append(list, subscriber{DID: row.DID, Rkey: row.Rkey, CreatedAt: row.CreatedAt}) + } + var count int64 if err := h.db.WithContext(c.Request().Context()).Model(&database.Subscription{}).Where("feed_id = ?", feedID).Count(&count).Error; err != nil { return writeError(c, http.StatusInternalServerError, "InternalError", "") @@ -51,6 +84,8 @@ func (h *Handlers) GetSubscribers(c echo.Context) error { return c.JSON(http.StatusOK, map[string]any{ "feedId": feedID, - "subscribers": count, + "subscribers": list, + "count": count, + "cursor": nextCursor, }) } diff --git a/appview/indexer/bookmark.go b/appview/indexer/bookmark.go index 1fac54a..cb27cd4 100644 --- a/appview/indexer/bookmark.go +++ b/appview/indexer/bookmark.go @@ -35,7 +35,10 @@ func (idx *Indexer) indexBookmark(ctx context.Context, did, rkey string, rec map } db := idx.db.WithContext(ctx) - return db.Where("did = ? AND rkey = ?", did, rkey).Assign(bookmark).FirstOrCreate(&bookmark).Error + if err := db.Where("did = ? AND rkey = ?", did, rkey).Assign(bookmark).FirstOrCreate(&bookmark).Error; err != nil { + return err + } + return idx.refreshEpisodeStats(ctx, feedID, episodeID) } func (idx *Indexer) indexProfile(ctx context.Context, did string, rec map[string]any) error { diff --git a/appview/indexer/comment.go b/appview/indexer/comment.go index fa7c0af..d49fab8 100644 --- a/appview/indexer/comment.go +++ b/appview/indexer/comment.go @@ -31,6 +31,7 @@ func (idx *Indexer) indexComment(ctx context.Context, did, rkey string, rec map[ PodcastGuid: asString(episode["podcastGuid"]), Text: asString(rec["text"]), CreatedAt: asString(rec["createdAt"]), + Facets: jsonBytes(rec["facets"]), } if ts, ok := asInt(rec["timestamp"]); ok { @@ -39,15 +40,19 @@ func (idx *Indexer) indexComment(ctx context.Context, did, rkey string, rec map[ if reply, ok := asMap(rec["reply"]); ok { if root, ok := asMap(reply["root"]); ok { - comment.RootURI = asString(root["uri"]) - comment.RootCID = asString(root["cid"]) + comment.ReplyRoot = asString(root["uri"]) } if parent, ok := asMap(reply["parent"]); ok { - comment.ParentURI = asString(parent["uri"]) - comment.ParentCID = asString(parent["cid"]) + comment.ReplyParent = asString(parent["uri"]) } } db := idx.db.WithContext(ctx) - return db.Where("did = ? AND rkey = ?", did, rkey).Assign(comment).FirstOrCreate(&comment).Error + if err := db.Where("did = ? AND rkey = ?", did, rkey).Assign(comment).FirstOrCreate(&comment).Error; err != nil { + return err + } + if err := idx.refreshPodcastStats(ctx, feedID); err != nil { + return err + } + return idx.refreshEpisodeStats(ctx, feedID, episodeID) } diff --git a/appview/indexer/indexer.go b/appview/indexer/indexer.go index b5ccaa4..ae96661 100644 --- a/appview/indexer/indexer.go +++ b/appview/indexer/indexer.go @@ -49,15 +49,49 @@ func (idx *Indexer) DeleteRecord(ctx context.Context, did, collection, rkey stri db := idx.db.WithContext(ctx) switch collection { case "xyz.effem.feed.subscription": - return db.Where("did = ? AND rkey = ?", did, rkey).Delete(&database.Subscription{}).Error + var row database.Subscription + if err := db.Where("did = ? AND rkey = ?", did, rkey).First(&row).Error; err != nil { + return nil + } + if err := db.Delete(&row).Error; err != nil { + return err + } + return idx.refreshPodcastStats(ctx, row.FeedID) case "xyz.effem.feed.comment": - return db.Where("did = ? AND rkey = ?", did, rkey).Delete(&database.Comment{}).Error + var row database.Comment + if err := db.Where("did = ? AND rkey = ?", did, rkey).First(&row).Error; err != nil { + return nil + } + if err := db.Delete(&row).Error; err != nil { + return err + } + if err := idx.refreshPodcastStats(ctx, row.FeedID); err != nil { + return err + } + return idx.refreshEpisodeStats(ctx, row.FeedID, row.EpisodeID) case "xyz.effem.feed.recommendation": - return db.Where("did = ? AND rkey = ?", did, rkey).Delete(&database.Recommendation{}).Error + var row database.Recommendation + if err := db.Where("did = ? AND rkey = ?", did, rkey).First(&row).Error; err != nil { + return nil + } + if err := db.Delete(&row).Error; err != nil { + return err + } + if err := idx.refreshPodcastStats(ctx, row.FeedID); err != nil { + return err + } + return idx.refreshEpisodeStats(ctx, row.FeedID, row.EpisodeID) case "xyz.effem.feed.list": return db.Where("did = ? AND rkey = ?", did, rkey).Delete(&database.PodcastList{}).Error case "xyz.effem.feed.bookmark": - return db.Where("did = ? AND rkey = ?", did, rkey).Delete(&database.Bookmark{}).Error + var row database.Bookmark + if err := db.Where("did = ? AND rkey = ?", did, rkey).First(&row).Error; err != nil { + return nil + } + if err := db.Delete(&row).Error; err != nil { + return err + } + return idx.refreshEpisodeStats(ctx, row.FeedID, row.EpisodeID) case "xyz.effem.actor.profile": return db.Where("did = ?", did).Delete(&database.Profile{}).Error default: diff --git a/appview/indexer/recommendation.go b/appview/indexer/recommendation.go index b9a27d4..2da2874 100644 --- a/appview/indexer/recommendation.go +++ b/appview/indexer/recommendation.go @@ -33,5 +33,11 @@ func (idx *Indexer) indexRecommendation(ctx context.Context, did, rkey string, r } db := idx.db.WithContext(ctx) - return db.Where("did = ? AND rkey = ?", did, rkey).Assign(reco).FirstOrCreate(&reco).Error + if err := db.Where("did = ? AND rkey = ?", did, rkey).Assign(reco).FirstOrCreate(&reco).Error; err != nil { + return err + } + if err := idx.refreshPodcastStats(ctx, feedID); err != nil { + return err + } + return idx.refreshEpisodeStats(ctx, feedID, episodeID) } diff --git a/appview/indexer/stats.go b/appview/indexer/stats.go new file mode 100644 index 0000000..ad2039c --- /dev/null +++ b/appview/indexer/stats.go @@ -0,0 +1,82 @@ +package indexer + +import ( + "context" + "time" + + "github.com/SparrowTek/effem-appview/appview/database" + "gorm.io/gorm/clause" +) + +func (idx *Indexer) refreshPodcastStats(ctx context.Context, feedID int) error { + if feedID <= 0 { + return nil + } + + db := idx.db.WithContext(ctx) + + var subCount int64 + if err := db.Model(&database.Subscription{}).Where("feed_id = ?", feedID).Count(&subCount).Error; err != nil { + return err + } + + var commentCount int64 + if err := db.Model(&database.Comment{}).Where("feed_id = ?", feedID).Count(&commentCount).Error; err != nil { + return err + } + + var recommendationCount int64 + if err := db.Model(&database.Recommendation{}).Where("feed_id = ?", feedID).Count(&recommendationCount).Error; err != nil { + return err + } + + row := database.PodcastStats{ + FeedID: feedID, + SubscriberCount: int(subCount), + CommentCount: int(commentCount), + RecommendationCount: int(recommendationCount), + LastUpdated: time.Now().UTC(), + } + + return db.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "feed_id"}}, + DoUpdates: clause.AssignmentColumns([]string{"subscriber_count", "comment_count", "recommendation_count", "last_updated"}), + }).Create(&row).Error +} + +func (idx *Indexer) refreshEpisodeStats(ctx context.Context, feedID, episodeID int) error { + if feedID <= 0 || episodeID <= 0 { + return nil + } + + db := idx.db.WithContext(ctx) + + var commentCount int64 + if err := db.Model(&database.Comment{}).Where("feed_id = ? AND episode_id = ?", feedID, episodeID).Count(&commentCount).Error; err != nil { + return err + } + + var recommendationCount int64 + if err := db.Model(&database.Recommendation{}).Where("feed_id = ? AND episode_id = ?", feedID, episodeID).Count(&recommendationCount).Error; err != nil { + return err + } + + var bookmarkCount int64 + if err := db.Model(&database.Bookmark{}).Where("feed_id = ? AND episode_id = ?", feedID, episodeID).Count(&bookmarkCount).Error; err != nil { + return err + } + + row := database.EpisodeStats{ + EpisodeID: episodeID, + FeedID: feedID, + CommentCount: int(commentCount), + RecommendationCount: int(recommendationCount), + BookmarkCount: int(bookmarkCount), + LastUpdated: time.Now().UTC(), + } + + return db.Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "episode_id"}}, + DoUpdates: clause.AssignmentColumns([]string{"feed_id", "comment_count", "recommendation_count", "bookmark_count", "last_updated"}), + }).Create(&row).Error +} diff --git a/appview/indexer/subscription.go b/appview/indexer/subscription.go index 39c2e84..fea2cc7 100644 --- a/appview/indexer/subscription.go +++ b/appview/indexer/subscription.go @@ -27,5 +27,8 @@ func (idx *Indexer) indexSubscription(ctx context.Context, did, rkey string, rec } db := idx.db.WithContext(ctx) - return db.Where("did = ? AND rkey = ?", did, rkey).Assign(sub).FirstOrCreate(&sub).Error + if err := db.Where("did = ? AND rkey = ?", did, rkey).Assign(sub).FirstOrCreate(&sub).Error; err != nil { + return err + } + return idx.refreshPodcastStats(ctx, feedID) } diff --git a/appview/podcastindex/cache.go b/appview/podcastindex/cache.go index c36898e..f3cc2b6 100644 --- a/appview/podcastindex/cache.go +++ b/appview/podcastindex/cache.go @@ -6,19 +6,11 @@ import ( "fmt" "time" + "github.com/SparrowTek/effem-appview/appview/database" "gorm.io/gorm" "gorm.io/gorm/clause" ) -type CachedResponse struct { - ID uint `gorm:"primaryKey"` - CacheKey string `gorm:"uniqueIndex;size:512"` - Response json.RawMessage `gorm:"type:jsonb"` - ExpiresAt time.Time `gorm:"index"` - CreatedAt time.Time `gorm:"autoCreateTime"` - UpdatedAt time.Time `gorm:"autoUpdateTime"` -} - type CachedClient struct { inner *Client db *gorm.DB @@ -28,19 +20,15 @@ func NewCachedClient(inner *Client, db *gorm.DB) *CachedClient { return &CachedClient{inner: inner, db: db} } -func RunMigrations(db *gorm.DB) error { - return db.AutoMigrate(&CachedResponse{}) -} - func (cc *CachedClient) getOrFetch( cacheKey string, ttl time.Duration, fetcher func() (json.RawMessage, error), ) (json.RawMessage, error) { - var cached CachedResponse + var cached database.PICache err := cc.db.Where("cache_key = ? AND expires_at > ?", cacheKey, time.Now()).First(&cached).Error if err == nil { - return cached.Response, nil + return json.RawMessage(cached.Response), nil } if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { return nil, err @@ -51,9 +39,9 @@ func (cc *CachedClient) getOrFetch( return nil, err } - row := CachedResponse{ + row := database.PICache{ CacheKey: cacheKey, - Response: data, + Response: []byte(data), ExpiresAt: time.Now().Add(ttl), } @@ -76,7 +64,7 @@ func (cc *CachedClient) SearchByTerm(query string, max int) (json.RawMessage, er func (cc *CachedClient) SearchEpisodesByTerm(query string, max int) (json.RawMessage, error) { key := fmt.Sprintf("search:episodes:%s:%d", query, max) - return cc.getOrFetch(key, 15*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch(key, time.Hour, func() (json.RawMessage, error) { return cc.inner.SearchEpisodesByTerm(query, max) }) } @@ -97,21 +85,21 @@ func (cc *CachedClient) GetEpisodesByFeedID(feedID int, max int) (json.RawMessag func (cc *CachedClient) GetEpisodeByID(episodeID int) (json.RawMessage, error) { key := fmt.Sprintf("episode:%d", episodeID) - return cc.getOrFetch(key, 15*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch(key, 6*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetEpisodeByID(episodeID) }) } func (cc *CachedClient) GetTrending(max int, lang string, categories string) (json.RawMessage, error) { key := fmt.Sprintf("trending:%d:%s:%s", max, lang, categories) - return cc.getOrFetch(key, 10*time.Minute, func() (json.RawMessage, error) { + return cc.getOrFetch(key, 30*time.Minute, func() (json.RawMessage, error) { return cc.inner.GetTrending(max, lang, categories) }) } func (cc *CachedClient) GetCategories() (json.RawMessage, error) { key := "categories:list" - return cc.getOrFetch(key, 24*time.Hour, func() (json.RawMessage, error) { + return cc.getOrFetch(key, 7*24*time.Hour, func() (json.RawMessage, error) { return cc.inner.GetCategories() }) } diff --git a/appview/server.go b/appview/server.go index 979c80d..16a9701 100644 --- a/appview/server.go +++ b/appview/server.go @@ -42,9 +42,6 @@ func NewServer(cfg Config) (*Server, error) { if err := database.RunMigrations(db); err != nil { return nil, fmt.Errorf("running database migrations: %w", err) } - if err := podcastindex.RunMigrations(db); err != nil { - return nil, fmt.Errorf("running podcast index cache migrations: %w", err) - } piClient := podcastindex.NewClient(cfg.PIKey, cfg.PISecret) cachedPI := podcastindex.NewCachedClient(piClient, db)