diff --git a/bookmark_handler.go b/bookmark_handler.go index 0f408db..4b4d6b8 100644 --- a/bookmark_handler.go +++ b/bookmark_handler.go @@ -9,6 +9,7 @@ import ( "net/http" "net/url" "strings" + "time" "github.com/bluesky-social/indigo/api/bsky" apibsky "github.com/bluesky-social/indigo/api/bsky" @@ -75,7 +76,7 @@ func (s *Server) HandleAddBookmark(w http.ResponseWriter, r *http.Request) { content = fmt.Sprintf("%s...", content[:75]) } - err = s.bookmarkStore.CreateBookmark(rkey, postURI, atPostURI, post.Author.Did, post.Author.Handle, usersDid, content) + err = s.bookmarkStore.CreateBookmark(rkey, postURI, atPostURI, post.Author.Did, post.Author.Handle, usersDid, content, time.Now().UnixMilli()) if err != nil { if errors.Is(err, store.ErrBookmarkAlreadyExists) { return @@ -131,10 +132,10 @@ func (s *Server) HandleDeleteBookmark(w http.ResponseWriter, r *http.Request) { return } - err = s.bookmarkStore.DeleteFeedPostsForBookmarkedPostURIandUserDID(bookmark.PostATURI, usersDid) + err = s.bookmarkStore.DeleteRepliedPostsForBookmarkedPostURIandUserDID(bookmark.PostATURI, usersDid) if err != nil { - slog.Error("deleting feed items for bookmark", "error", err) - http.Error(w, "deleting feed items for bookmark", http.StatusInternalServerError) + slog.Error("deleting replied posts for bookmark", "error", err) + http.Error(w, "deleting replied posts for bookmark", http.StatusInternalServerError) return } diff --git a/dm_handler.go b/dm_handler.go index 07db58b..197d56a 100644 --- a/dm_handler.go +++ b/dm_handler.go @@ -243,7 +243,7 @@ func (d *DmService) handleCreateBookmark(msg Message) error { rkey := getRKeyFromATURI(msg.Embed.Record.URI) - err := d.bookmarkStore.CreateBookmark(rkey, publicURI, msg.Embed.Record.URI, msg.Embed.Record.Author.Did, msg.Embed.Record.Author.Handle, msg.Sender.Did, content) + err := d.bookmarkStore.CreateBookmark(rkey, publicURI, msg.Embed.Record.URI, msg.Embed.Record.Author.Did, msg.Embed.Record.Author.Handle, msg.Sender.Did, content, time.Now().UnixMilli()) if err != nil { return fmt.Errorf("creating bookmark: %w", err) } @@ -253,9 +253,9 @@ func (d *DmService) handleCreateBookmark(msg Message) error { func (d *DmService) handleDeleteBookmark(msg Message) error { rkey := getRKeyFromATURI(msg.Embed.Record.URI) - err := d.bookmarkStore.DeleteFeedPostsForBookmarkedPostURIandUserDID(msg.Embed.Record.URI, msg.Sender.Did) + err := d.bookmarkStore.DeleteRepliedPostsForBookmarkedPostURIandUserDID(msg.Embed.Record.URI, msg.Sender.Did) if err != nil { - return fmt.Errorf("failed to delete feed posts of replies to bookmark for user: %w", err) + return fmt.Errorf("failed to delete replied posts for bookmark for user: %w", err) } err = d.bookmarkStore.DeleteBookmark(rkey, msg.Sender.Did) diff --git a/feed_handlers.go b/feed_handlers.go index 1a76af9..0865295 100644 --- a/feed_handlers.go +++ b/feed_handlers.go @@ -104,7 +104,10 @@ func (s *Server) HandleDescribeFeedGenerator(w http.ResponseWriter, r *http.Requ DID: fmt.Sprintf("did:web:%s", s.feedHost), Feeds: []FeedRespsonse{ { - URI: fmt.Sprintf("at://%s/app.bsky.feed.generator/wills-test", s.feedDidBase), + URI: fmt.Sprintf("at://%s/app.bsky.feed.generator/bookmark-replies", s.feedDidBase), + }, + { + URI: fmt.Sprintf("at://%s/app.bsky.feed.generator/bookmarks", s.feedDidBase), }, }, } diff --git a/feedgenerator.go b/feedgenerator.go index da5bf45..250ad09 100644 --- a/feedgenerator.go +++ b/feedgenerator.go @@ -5,26 +5,42 @@ import ( "fmt" "log/slog" "strconv" + "strings" "github.com/willdot/bskyfeedgen/store" ) -type feedStore interface { - GetUsersFeed(usersDID string, cursor int64, limit int) ([]store.FeedPost, error) - AddFeedPost(feedPost store.FeedPost) error +type repliesStore interface { + GetUsersReplies(usersDID string, cursor int64, limit int) ([]store.ReplyPost, error) + GetBookmarksForUserWithPaging(userDID string, cursor int64, limit int) ([]store.Bookmark, error) + AddRepliedPost(replyPost store.ReplyPost) error } type FeedGenerator struct { - store feedStore + store repliesStore } -func NewFeedGenerator(store feedStore) *FeedGenerator { +func NewFeedGenerator(store repliesStore) *FeedGenerator { return &FeedGenerator{ store: store, } } func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor string, limit int) (FeedReponse, error) { + switch { + case strings.Contains(feed, "bookmark-replies"): + return f.getBookmarkRepliesFeed(ctx, userDID, cursor, limit) + case strings.Contains(feed, "bookmarks"): + return f.getBookmarksFeed(ctx, userDID, cursor, limit) + + default: + return FeedReponse{ + Feed: make([]FeedItem, 0), + }, fmt.Errorf("invalid feed requested") + } +} + +func (f *FeedGenerator) getBookmarkRepliesFeed(ctx context.Context, userDID, cursor string, limit int) (FeedReponse, error) { resp := FeedReponse{ Feed: make([]FeedItem, 0), } @@ -38,13 +54,13 @@ func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor strin cursorInt = 9999999999999 } - usersFeed, err := f.store.GetUsersFeed(userDID, int64(cursorInt), limit) + usersReplies, err := f.store.GetUsersReplies(userDID, int64(cursorInt), limit) if err != nil { - return resp, fmt.Errorf("get users feed items from DB: %w", err) + return resp, fmt.Errorf("get users replies from DB: %w", err) } - feedItems := make([]FeedItem, 0, len(usersFeed)) - for _, post := range usersFeed { + feedItems := make([]FeedItem, 0, len(usersReplies)) + for _, post := range usersReplies { feedItems = append(feedItems, FeedItem{ Post: post.ReplyURI, }) @@ -54,8 +70,45 @@ func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor strin // only set the return cursor if there was a record returned and that the len of records // being returned is the same as the limit - if len(usersFeed) > 0 && len(usersFeed) == limit { - lastFeedItem := usersFeed[len(usersFeed)-1] + if len(usersReplies) > 0 && len(usersReplies) == limit { + lastFeedItem := usersReplies[len(usersReplies)-1] + resp.Cursor = fmt.Sprintf("%d", lastFeedItem.CreatedAt) + } + return resp, nil +} + +func (f *FeedGenerator) getBookmarksFeed(ctx context.Context, userDID, cursor string, limit int) (FeedReponse, error) { + resp := FeedReponse{ + Feed: make([]FeedItem, 0), + } + + cursorInt, err := strconv.Atoi(cursor) + if err != nil && cursor != "" { + slog.Error("convert cursor to int", "error", err, "cursor value", cursor) + } + if cursorInt == 0 { + // if no cursor provided use a date waaaaay in the future to start the less than query + cursorInt = 9999999999999 + } + + usersBookmarks, err := f.store.GetBookmarksForUserWithPaging(userDID, int64(cursorInt), limit) + if err != nil { + return resp, fmt.Errorf("get users bookmarks from DB: %w", err) + } + + feedItems := make([]FeedItem, 0, len(usersBookmarks)) + for _, bookmark := range usersBookmarks { + feedItems = append(feedItems, FeedItem{ + Post: bookmark.PostATURI, + }) + } + + resp.Feed = feedItems + + // only set the return cursor if there was a record returned and that the len of records + // being returned is the same as the limit + if len(usersBookmarks) > 0 && len(usersBookmarks) == limit { + lastFeedItem := usersBookmarks[len(usersBookmarks)-1] resp.Cursor = fmt.Sprintf("%d", lastFeedItem.CreatedAt) } return resp, nil diff --git a/firehose_handler.go b/firehose_handler.go index ac357a3..465887e 100644 --- a/firehose_handler.go +++ b/firehose_handler.go @@ -14,7 +14,7 @@ import ( ) type HandlerStore interface { - AddFeedPost(feedItem store.FeedPost) error + AddRepliedPost(replyPost store.ReplyPost) error GetBookmarksForPost(postURI string) ([]string, error) } @@ -68,7 +68,7 @@ func (h *handler) handleCreateEvent(_ context.Context, event *models.Event) erro } replyPostURI := fmt.Sprintf("at://%s/app.bsky.feed.post/%s", event.Did, event.Commit.RKey) - h.createFeedPostForSubscribedUsers(subscribedDids, replyPostURI, subscribedPostURI, createdAt.UnixMilli()) + h.createReplyPostForSubscribedUsers(subscribedDids, replyPostURI, subscribedPostURI, createdAt.UnixMilli()) return nil } @@ -83,17 +83,17 @@ func (h *handler) getSubscribedDidsForPost(postURI string) []string { return dids } -func (h *handler) createFeedPostForSubscribedUsers(usersDids []string, replyPostURI, subscribedPostURI string, createdAt int64) { +func (h *handler) createReplyPostForSubscribedUsers(usersDids []string, replyPostURI, subscribedPostURI string, createdAt int64) { for _, did := range usersDids { - feedItem := store.FeedPost{ + repliedPost := store.ReplyPost{ ReplyURI: replyPostURI, UserDID: did, SubscribedPostURI: subscribedPostURI, CreatedAt: createdAt, } - err := h.store.AddFeedPost(feedItem) + err := h.store.AddRepliedPost(repliedPost) if err != nil { - slog.Error("add users feed item", "error", err, "did", did, "reply post URI", replyPostURI) + slog.Error("add users replied post", "error", err, "did", did, "reply post URI", replyPostURI) _ = bugsnag.Notify(err) continue } diff --git a/server.go b/server.go index c7b2353..14cae0a 100644 --- a/server.go +++ b/server.go @@ -28,11 +28,11 @@ type Store interface { } type BookmarkStore interface { - CreateBookmark(postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content string) error + CreateBookmark(postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content string, createdAt int64) error GetBookmarksForUser(userDID string) ([]store.Bookmark, error) DeleteBookmark(postRKey, userDID string) error GetBookmarkByRKeyForUser(rkey, userDID string) (*store.Bookmark, error) - DeleteFeedPostsForBookmarkedPostURIandUserDID(subscribedPostURI, userDID string) error + DeleteRepliedPostsForBookmarkedPostURIandUserDID(subscribedPostURI, userDID string) error } type OauthRequestStore interface { diff --git a/store/bookmark.go b/store/bookmark.go index 75e00d0..423663b 100644 --- a/store/bookmark.go +++ b/store/bookmark.go @@ -19,6 +19,7 @@ func createBookmarksTable(db *sql.DB) error { "authorHandle" TEXT, "userDID" TEXT, "content" TEXT, + "createdAt" integer NOT NULL, UNIQUE(postRKey, userDID) );` @@ -45,11 +46,12 @@ type Bookmark struct { AuthorHandle string UserDID string Content string + CreatedAt int64 } -func (s *Store) CreateBookmark(postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content string) error { - sql := `INSERT INTO bookmarks (postRKey, postURI,postATURI, authorDID, authorHandle, userDID, content) VALUES (?, ?, ?, ?, ?, ?, ?) ON CONFLICT(postRKey, userDID) DO NOTHING;` - res, err := s.db.Exec(sql, postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content) +func (s *Store) CreateBookmark(postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content string, createdAt int64) error { + sql := `INSERT INTO bookmarks (postRKey, postURI,postATURI, authorDID, authorHandle, userDID, content, createdAt) VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(postRKey, userDID) DO NOTHING;` + res, err := s.db.Exec(sql, postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content, createdAt) if err != nil { return fmt.Errorf("exec insert bookmark: %w", err) } @@ -61,7 +63,7 @@ func (s *Store) CreateBookmark(postRKey, postURI, postATURI, authorDID, authorHa } func (s *Store) GetBookmarksForUser(userDID string) ([]Bookmark, error) { - sql := "SELECT id, postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content FROM bookmarks WHERE userDID = ?;" + sql := "SELECT id, postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content, createdAt FROM bookmarks WHERE userDID = ?;" rows, err := s.db.Query(sql, userDID) if err != nil { return nil, fmt.Errorf("run query to get bookmarked posts for user: %w", err) @@ -71,7 +73,29 @@ func (s *Store) GetBookmarksForUser(userDID string) ([]Bookmark, error) { var results []Bookmark for rows.Next() { var bookmark Bookmark - if err := rows.Scan(&bookmark.ID, &bookmark.PostRKey, &bookmark.PostURI, &bookmark.PostATURI, &bookmark.AuthorDID, &bookmark.AuthorHandle, &bookmark.UserDID, &bookmark.Content); err != nil { + if err := rows.Scan(&bookmark.ID, &bookmark.PostRKey, &bookmark.PostURI, &bookmark.PostATURI, &bookmark.AuthorDID, &bookmark.AuthorHandle, &bookmark.UserDID, &bookmark.Content, &bookmark.CreatedAt); err != nil { + return nil, fmt.Errorf("scan row: %w", err) + } + + results = append(results, bookmark) + } + return results, nil +} + +func (s *Store) GetBookmarksForUserWithPaging(userDID string, cursor int64, limit int) ([]Bookmark, error) { + sql := `SELECT id, postRKey, postURI, postATURI, authorDID, authorHandle, userDID, content, createdAt FROM bookmarks + WHERE userDID = ? AND createdAt < ? + ORDER BY createdAt DESC LIMIT ?;` + rows, err := s.db.Query(sql, userDID, cursor, limit) + if err != nil { + return nil, fmt.Errorf("run query to get bookmarked posts for user: %w", err) + } + defer rows.Close() + + var results []Bookmark + for rows.Next() { + var bookmark Bookmark + if err := rows.Scan(&bookmark.ID, &bookmark.PostRKey, &bookmark.PostURI, &bookmark.PostATURI, &bookmark.AuthorDID, &bookmark.AuthorHandle, &bookmark.UserDID, &bookmark.Content, &bookmark.CreatedAt); err != nil { return nil, fmt.Errorf("scan row: %w", err) } diff --git a/store/database.go b/store/database.go index 0dd733e..8a7c608 100644 --- a/store/database.go +++ b/store/database.go @@ -32,9 +32,9 @@ func New(dbPath string) (*Store, error) { return nil, fmt.Errorf("ping db: %w", err) } - err = createFeedTable(db) + err = createRepliesTable(db) if err != nil { - return nil, fmt.Errorf("creating feed table: %w", err) + return nil, fmt.Errorf("creating replies table: %w", err) } err = createBookmarksTable(db) diff --git a/store/feed.go b/store/feed.go deleted file mode 100644 index 8d350c9..0000000 --- a/store/feed.go +++ /dev/null @@ -1,87 +0,0 @@ -package store - -import ( - "database/sql" - "fmt" - "log/slog" -) - -func createFeedTable(db *sql.DB) error { - createFeedTableSQL := `CREATE TABLE IF NOT EXISTS feed ( - "id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, - "replyURI" TEXT, - "userDID" TEXT, - "subscribedPostURI" TEXT, - "createdAt" integer NOT NULL, - UNIQUE(replyURI, userDID) - );` - - slog.Info("Create feed table...") - statement, err := db.Prepare(createFeedTableSQL) - if err != nil { - return fmt.Errorf("prepare DB statement to create feeds table: %w", err) - } - _, err = statement.Exec() - if err != nil { - return fmt.Errorf("exec sql statement to create feeds table: %w", err) - } - slog.Info("feed table created") - - return nil -} - -type FeedPost struct { - ID int - ReplyURI string - UserDID string - SubscribedPostURI string - CreatedAt int64 -} - -func (s *Store) AddFeedPost(feedPost FeedPost) error { - sql := `INSERT INTO feed (replyURI, userDID, subscribedPostURI, createdAt) VALUES (?, ?, ?, ?) ON CONFLICT(replyURI, userDID) DO NOTHING;` - _, err := s.db.Exec(sql, feedPost.ReplyURI, feedPost.UserDID, feedPost.SubscribedPostURI, feedPost.CreatedAt) - if err != nil { - return fmt.Errorf("exec insert feed item: %w", err) - } - return nil -} - -func (s *Store) GetUsersFeed(usersDID string, cursor int64, limit int) ([]FeedPost, error) { - sql := `SELECT id, replyURI, userDID, subscribedPostURI, createdAt FROM feed - WHERE userDID = ? AND createdAt < ? - ORDER BY createdAt DESC LIMIT ?;` - rows, err := s.db.Query(sql, usersDID, cursor, limit) - if err != nil { - return nil, fmt.Errorf("run query to get users feed posts: %w", err) - } - defer rows.Close() - - feedPosts := make([]FeedPost, 0) - for rows.Next() { - var feedPost FeedPost - if err := rows.Scan(&feedPost.ID, &feedPost.ReplyURI, &feedPost.UserDID, &feedPost.SubscribedPostURI, &feedPost.CreatedAt); err != nil { - return nil, fmt.Errorf("scan row: %w", err) - } - feedPosts = append(feedPosts, feedPost) - } - - return feedPosts, nil -} - -func (s *Store) DeleteFeedPostsForBookmarkedPostURIandUserDID(subscribedPostURI, userDID string) error { - sql := "DELETE FROM feed WHERE subscribedPostURI = ? AND userDID = ?;" - statement, err := s.db.Prepare(sql) - if err != nil { - return fmt.Errorf("prepare delete feed posts: %w", err) - } - res, err := statement.Exec(subscribedPostURI, userDID) - if err != nil { - return fmt.Errorf("exec delete feed posts: %w", err) - } - - n, _ := res.RowsAffected() - - slog.Info("delete feed posts result", "affected rows", n) - return nil -} diff --git a/store/replies.go b/store/replies.go new file mode 100644 index 0000000..2396b01 --- /dev/null +++ b/store/replies.go @@ -0,0 +1,87 @@ +package store + +import ( + "database/sql" + "fmt" + "log/slog" +) + +func createRepliesTable(db *sql.DB) error { + createRepliesTableSQL := `CREATE TABLE IF NOT EXISTS replies ( + "id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, + "replyURI" TEXT, + "userDID" TEXT, + "subscribedPostURI" TEXT, + "createdAt" integer NOT NULL, + UNIQUE(replyURI, userDID) + );` + + slog.Info("Create replies table...") + statement, err := db.Prepare(createRepliesTableSQL) + if err != nil { + return fmt.Errorf("prepare DB statement to create replies table: %w", err) + } + _, err = statement.Exec() + if err != nil { + return fmt.Errorf("exec sql statement to create replies table: %w", err) + } + slog.Info("replies table created") + + return nil +} + +type ReplyPost struct { + ID int + ReplyURI string + UserDID string + SubscribedPostURI string + CreatedAt int64 +} + +func (s *Store) AddRepliedPost(replyPost ReplyPost) error { + sql := `INSERT INTO replies (replyURI, userDID, subscribedPostURI, createdAt) VALUES (?, ?, ?, ?) ON CONFLICT(replyURI, userDID) DO NOTHING;` + _, err := s.db.Exec(sql, replyPost.ReplyURI, replyPost.UserDID, replyPost.SubscribedPostURI, replyPost.CreatedAt) + if err != nil { + return fmt.Errorf("exec insert replies post: %w", err) + } + return nil +} + +func (s *Store) GetUsersReplies(usersDID string, cursor int64, limit int) ([]ReplyPost, error) { + sql := `SELECT id, replyURI, userDID, subscribedPostURI, createdAt FROM replies + WHERE userDID = ? AND createdAt < ? + ORDER BY createdAt DESC LIMIT ?;` + rows, err := s.db.Query(sql, usersDID, cursor, limit) + if err != nil { + return nil, fmt.Errorf("run query to get users replied posts: %w", err) + } + defer rows.Close() + + repliedPosts := make([]ReplyPost, 0) + for rows.Next() { + var replyPost ReplyPost + if err := rows.Scan(&replyPost.ID, &replyPost.ReplyURI, &replyPost.UserDID, &replyPost.SubscribedPostURI, &replyPost.CreatedAt); err != nil { + return nil, fmt.Errorf("scan row: %w", err) + } + repliedPosts = append(repliedPosts, replyPost) + } + + return repliedPosts, nil +} + +func (s *Store) DeleteRepliedPostsForBookmarkedPostURIandUserDID(subscribedPostURI, userDID string) error { + sql := "DELETE FROM replies WHERE subscribedPostURI = ? AND userDID = ?;" + statement, err := s.db.Prepare(sql) + if err != nil { + return fmt.Errorf("prepare delete replies: %w", err) + } + res, err := statement.Exec(subscribedPostURI, userDID) + if err != nil { + return fmt.Errorf("exec delete replies: %w", err) + } + + n, _ := res.RowsAffected() + + slog.Info("delete replies result", "affected rows", n) + return nil +}