diff --git a/OLDDockerfile b/OLDDockerfile deleted file mode 100644 index 6230ff7..0000000 --- a/OLDDockerfile +++ /dev/null @@ -1,21 +0,0 @@ -# # Use the Go 1.23 alpine official image -# # https://hub.docker.com/_/golang -# FROM golang:1.23-alpine - -# # Create and change to the app directory. -# WORKDIR /app - -# # Copy go mod and sum files -# COPY go.mod go.sum ./ - -# # Copy local code to the container image. -# COPY . ./ - -# # Install project dependencies -# RUN CGO_ENABLED=1 go mod download - -# # Build the app -# RUN go build -o app - -# # Run the service on container startup. -# ENTRYPOINT ["./app"] diff --git a/consumer.go b/consumer.go index 36b916c..745fc23 100644 --- a/consumer.go +++ b/consumer.go @@ -2,7 +2,7 @@ package main import ( "context" - "database/sql" + "encoding/json" "fmt" "log/slog" @@ -14,6 +14,7 @@ import ( "github.com/bluesky-social/jetstream/pkg/client/schedulers/sequential" "github.com/bluesky-social/jetstream/pkg/models" "github.com/bugsnag/bugsnag-go/v2" + "github.com/willdot/bskyfeedgen/store" ) type consumer struct { @@ -38,7 +39,7 @@ func (con *consumer) Consume(ctx context.Context, feedGen *FeedGenerator, logger h := &handler{ seenSeqs: make(map[int64]struct{}), feedGenerator: feedGen, - db: feedGen.db, + store: *feedGen.store, } scheduler := sequential.NewScheduler("jetstream_localdev", logger, h.HandleEvent) @@ -64,7 +65,7 @@ type handler struct { seenSeqs map[int64]struct{} highwater int64 feedGenerator *FeedGenerator - db *sql.DB + store store.Store } func (h *handler) HandleEvent(ctx context.Context, event *models.Event) error { @@ -98,23 +99,24 @@ func (h *handler) handleCreateEvent(_ context.Context, event *models.Event) erro return nil } - parentURI := post.Reply.Parent.Uri + subscribedPostURI := post.Reply.Parent.Uri - // look for posts where I've "subsribed" so that we can add the parent URI to a list of replies to that parent to look for + // look for posts that are "subscribe" so that we can add the post URI to a list of posts we want to find replies for if strings.Contains(post.Text, "/subscribe") && event.Did == "did:plc:dadhhalkfcq3gucaq25hjqon" { - slog.Info("a post that's subscribing to a parent. Adding to parents to look for", "parent URI", parentURI) - return h.addDidToSubscribedParent(parentURI, event.Did, event.Commit.RKey) + slog.Info("a post that's subscribing to another post. Adding to posts to look for", "subscribed post URI", subscribedPostURI) + return h.addDidToSubscribedPost(subscribedPostURI, event.Did, event.Commit.RKey) } // see if the post is a reply to a post we are subscribed to - subscribedDids := h.getSubscribedDidsForParent(parentURI) + subscribedDids := h.getSubscribedDidsForPost(subscribedPostURI) if len(subscribedDids) == 0 { return nil } - slog.Info("post is a reply to a parent that users are subscribed to", "parent URI", parentURI, "dids", subscribedDids, "RKey", event.Commit.RKey) + slog.Info("post is a reply to a post that users are subscribed to", "subscribed post URI", subscribedPostURI, "dids", subscribedDids, "RKey", event.Commit.RKey) - h.feedGenerator.AddToFeedPosts(subscribedDids, parentURI, fmt.Sprintf("at://%s/app.bsky.feed.post/%s", event.Did, event.Commit.RKey)) + replyPostURI := fmt.Sprintf("at://%s/app.bsky.feed.post/%s", event.Did, event.Commit.RKey) + h.feedGenerator.AddToFeedPosts(subscribedDids, subscribedPostURI, replyPostURI) return nil } @@ -128,45 +130,42 @@ func (h *handler) handleDeleteEvent(_ context.Context, event *models.Event) erro return nil } slog.Info("delete event received", "did", event.Did, "rkey", event.Commit.RKey) - - parentURI, err := getSubscribingPostParentURI(h.db, event.Did, event.Commit.RKey) + subscribedPostURI, err := h.store.GetSubscribedPostURI(event.Did, event.Commit.RKey) if err != nil { - slog.Error("get subscribing post parent URI", "error", err, "rkey", event.Commit.RKey, "user DID", event.Did) - return fmt.Errorf("get subscribing post parent URI: %w", err) + slog.Error("get subscribed post URI", "error", err, "rkey", event.Commit.RKey, "user DID", event.Did) + return fmt.Errorf("get subscribed post URI: %w", err) } - slog.Info("delete parent URI", "parent URI", parentURI, "rkey", event.Commit.RKey) - - // delete from feeds for the parentURI and the users DID first. This is so that if this fails, it can be tried again and the + // delete from feeds for the subscribedPostURI and the users DID first. This is so that if this fails, it can be tried again and the // subscription will be still there - err = deleteFeedItemsForParentURIandUserDID(h.db, parentURI, event.Did) + err = h.store.DeleteFeedItemsForSubscribedPostURIandUserDID(subscribedPostURI, event.Did) if err != nil { - slog.Error("delete feed items for parentURI and user", "error", err, "parentURI", parentURI, "user DID", event.Did) - return fmt.Errorf("delete feed items for parentURI and user: %w", err) + slog.Error("delete feed items for subscribedPostURI and user", "error", err, "subscribedPostURI", subscribedPostURI, "user DID", event.Did) + return fmt.Errorf("delete feed items for subscribedPostURI and user: %w", err) } - // delete from subscriptions for the parentURI and the users DID now that we have cleaned up the feeds - err = deleteSubscriptionForUser(h.db, event.Did, parentURI) + // delete from subscriptions for the postURI and the users DID now that we have cleaned up the feeds + err = h.store.DeleteSubscriptionForUser(event.Did, subscribedPostURI) if err != nil { - slog.Error("delete subscription for user", "error", err, "parentURI", parentURI, "user DID", event.Did) + slog.Error("delete subscription for user", "error", err, "subscribedPostURI", subscribedPostURI, "user DID", event.Did) return fmt.Errorf("delete subscription and user: %w", err) } return nil } -func (h *handler) addDidToSubscribedParent(parentURI, userDid, rkey string) error { - err := addSubscriptionForParent(h.db, parentURI, userDid, rkey) +func (h *handler) addDidToSubscribedPost(subscribedPostURI, userDid, rkey string) error { + err := h.store.AddSubscriptionForPost(subscribedPostURI, userDid, rkey) if err != nil { - return fmt.Errorf("add subscription for parent: %w", err) + return fmt.Errorf("add subscription for post: %w", err) } return nil } -func (h *handler) getSubscribedDidsForParent(parentURI string) []string { - dids, err := getSubscriptionsForParent(h.db, parentURI) +func (h *handler) getSubscribedDidsForPost(postURI string) []string { + dids, err := h.store.GetSubscriptionsForPost(postURI) if err != nil { - slog.Error("getting subscriptions for parent", "error", err) + slog.Error("getting subscriptions for post", "error", err) bugsnag.Notify(err) } diff --git a/database.go b/database.go deleted file mode 100644 index 71f2c5c..0000000 --- a/database.go +++ /dev/null @@ -1,226 +0,0 @@ -package main - -import ( - "context" - "database/sql" - "errors" - "fmt" - "log/slog" - "os" - - _ "github.com/glebarez/go-sqlite" -) - -func NewDatabase(dbPath string) (*sql.DB, error) { - err := createDbFile(dbPath) - if err != nil { - return nil, fmt.Errorf("create db file: %w", err) - } - - db, err := sql.Open("sqlite", dbPath) - if err != nil { - return nil, fmt.Errorf("open database: %w", err) - } - - err = db.Ping() - if err != nil { - return nil, fmt.Errorf("ping db: %w", err) - } - - err = createFeedTable(db) - if err != nil { - return nil, fmt.Errorf("creating feed table: %w", err) - } - - err = createSubscriptionsTable(db) - if err != nil { - return nil, fmt.Errorf("creating subscription table: %w", err) - } - - return db, nil -} - -func createDbFile(dbFilename string) error { - if _, err := os.Stat(dbFilename); !errors.Is(err, os.ErrNotExist) { - return nil - } - - f, err := os.Create(dbFilename) - if err != nil { - return fmt.Errorf("create db file : %w", err) - } - f.Close() - return nil -} - -func createFeedTable(db *sql.DB) error { - createFeedTableSQL := `CREATE TABLE IF NOT EXISTS feed ( - "id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, - "uri" TEXT, - "userDID" TEXT, - "parentURI" TEXT, - UNIQUE(uri, 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 -} - -func createSubscriptionsTable(db *sql.DB) error { - createSubscriptionsTableSQL := `CREATE TABLE IF NOT EXISTS subscriptions ( - "id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, - "parentURI" TEXT, - "userDID" TEXT, - "subscriptionRkey" TEXT, - UNIQUE(parentURI, userDID) - );` - - slog.Info("Create subscriptions table...") - statement, err := db.Prepare(createSubscriptionsTableSQL) - if err != nil { - return fmt.Errorf("prepare DB statement to create subscriptions table: %w", err) - } - _, err = statement.Exec() - if err != nil { - return fmt.Errorf("exec sql statement to create subscriptions table: %w", err) - } - slog.Info("subscriptions table created") - - return nil -} - -type feedItem struct { - ID int - URI string - UserDID string - parentURI string -} - -func addFeedItem(_ context.Context, db *sql.DB, feedItem feedItem) error { - slog.Info("add feed item", "parenturi", feedItem.parentURI, "user did", feedItem.UserDID) - sql := `INSERT INTO feed (uri, userDID, parentURI) VALUES (?, ?, ?) ON CONFLICT(uri, userDID) DO NOTHING;` - _, err := db.Exec(sql, feedItem.URI, feedItem.UserDID, feedItem.parentURI) - if err != nil { - return fmt.Errorf("exec insert feed item: %w", err) - } - return nil -} - -func getUsersFeedItems(db *sql.DB, usersDID string) ([]feedItem, error) { - sql := "SELECT id, uri, userDID FROM feed WHERE userDID = ?;" - rows, err := db.Query(sql, usersDID) - if err != nil { - return nil, fmt.Errorf("run query to get users feed item: %w", err) - } - defer rows.Close() - - feedItems := make([]feedItem, 0) - for rows.Next() { - var feedItem feedItem - if err := rows.Scan(&feedItem.ID, &feedItem.URI, &feedItem.UserDID); err != nil { - return nil, fmt.Errorf("scan row: %w", err) - } - feedItems = append(feedItems, feedItem) - } - - return feedItems, nil -} - -func deleteFeedItemsForParentURIandUserDID(db *sql.DB, parentURI, userDID string) error { - slog.Info("delete feed", "parent uri", parentURI, "userdid", userDID) - - sql := "DELETE FROM feed WHERE parentURI = ? AND userDID = ?;" - statement, err := db.Prepare(sql) - if err != nil { - return fmt.Errorf("prepare delete feed items: %w", err) - } - res, err := statement.Exec(parentURI, userDID) - if err != nil { - return fmt.Errorf("exec delete feed items: %w", err) - } - - n, _ := res.RowsAffected() - - slog.Info("delete feed res", "affected rows", n) - return nil -} - -type subscription struct { - ID int - ParentURI string - UserDID string - SubecriptionRkey string -} - -func getSubscriptionsForParent(db *sql.DB, parentURI string) ([]string, error) { - sql := "SELECT id, parentURI, userDID FROM subscriptions WHERE parentURI = ?" - rows, err := db.Query(sql, parentURI) - if err != nil { - return nil, fmt.Errorf("run query to get subscriptions: %w", err) - } - defer rows.Close() - - dids := make([]string, 0) - for rows.Next() { - var subscription subscription - if err := rows.Scan(&subscription.ID, &subscription.ParentURI, &subscription.UserDID); err != nil { - return nil, fmt.Errorf("scan row: %w", err) - } - dids = append(dids, subscription.UserDID) - } - - return dids, nil -} - -// urh -func addSubscriptionForParent(db *sql.DB, parentURI, userDid, subscriptionRkey string) error { - sql := `INSERT INTO subscriptions (parentURI, userDID, subscriptionRkey) VALUES (?, ?, ?) ON CONFLICT(parentURI, userDID) DO NOTHING;` - _, err := db.Exec(sql, parentURI, userDid, subscriptionRkey) - if err != nil { - return fmt.Errorf("exec insert subscrptions: %w", err) - } - return nil -} - -func getSubscribingPostParentURI(db *sql.DB, userDID, rkey string) (string, error) { - slog.Info("params", "rkey", rkey, "did", userDID) - sql := "SELECT id, parentURI FROM subscriptions WHERE subscriptionRkey = ? AND userDID = ?;" - rows, err := db.Query(sql, rkey, userDID) - if err != nil { - return "", fmt.Errorf("run query to get subscribing post parent URI: %w", err) - } - defer rows.Close() - - parentURI := "" - for rows.Next() { - var subscription subscription - if err := rows.Scan(&subscription.ID, &subscription.ParentURI); err != nil { - return "", fmt.Errorf("scan row: %w", err) - } - - slog.Info("record", "val", subscription) - - parentURI = subscription.ParentURI - break - } - return parentURI, nil -} - -func deleteSubscriptionForUser(db *sql.DB, userDID, parentURI string) error { - sql := "DELETE FROM subscriptions WHERE parentURI = ? AND userDID = ?;" - _, err := db.Exec(sql, parentURI, userDID) - if err != nil { - return fmt.Errorf("exec delete subscription for user: %w", err) - } - return nil -} diff --git a/feed.go b/feed.go index 8b33b36..257aa83 100644 --- a/feed.go +++ b/feed.go @@ -2,20 +2,20 @@ package main import ( "context" - "database/sql" "fmt" "log/slog" "github.com/bugsnag/bugsnag-go/v2" + "github.com/willdot/bskyfeedgen/store" ) type FeedGenerator struct { - db *sql.DB + store *store.Store } -func NewFeedGenerator(db *sql.DB) *FeedGenerator { +func NewFeedGenerator(store *store.Store) *FeedGenerator { return &FeedGenerator{ - db: db, + store: store, } } @@ -24,7 +24,7 @@ func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor strin Feed: make([]FeedItem, 0, 0), } - usersFeed, err := getUsersFeedItems(f.db, userDID) + usersFeed, err := f.store.GetUsersFeedItems(userDID) if err != nil { return resp, fmt.Errorf("get users feed items from DB: %w", err) } @@ -32,7 +32,7 @@ func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor strin feedItems := make([]FeedItem, 0, len(usersFeed)) for _, post := range usersFeed { feedItems = append(feedItems, FeedItem{ - Post: post.URI, + Post: post.ReplyURI, }) } @@ -42,16 +42,16 @@ func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor strin return resp, nil } -func (f *FeedGenerator) AddToFeedPosts(usersDids []string, parentURI, postURI string) { +func (f *FeedGenerator) AddToFeedPosts(usersDids []string, subscribedPostURI, replyPostURI string) { for _, did := range usersDids { - feedItem := feedItem{ - URI: postURI, - UserDID: did, - parentURI: parentURI, + feedItem := store.FeedItem{ + ReplyURI: replyPostURI, + UserDID: did, + SubscribedPostURI: subscribedPostURI, } - err := addFeedItem(context.Background(), f.db, feedItem) + err := f.store.AddFeedItem(feedItem) if err != nil { - slog.Error("add users feed item", "error", err, "did", did, "uri", postURI) + slog.Error("add users feed item", "error", err, "did", did, "reply post URI", replyPostURI) bugsnag.Notify(err) continue } diff --git a/main.go b/main.go index b4f7379..c429177 100644 --- a/main.go +++ b/main.go @@ -13,6 +13,7 @@ import ( "github.com/avast/retry-go/v4" "github.com/bugsnag/bugsnag-go/v2" + "github.com/willdot/bskyfeedgen/store" ) const ( @@ -44,15 +45,16 @@ func main() { return } dbFilename := path.Join(dbMountPath, "database.db") - db, err := NewDatabase(dbFilename) + + store, err := store.New(dbFilename) if err != nil { - slog.Error("create new database", "error", err) + slog.Error("create new store", "error", err) bugsnag.Notify(err) return } - defer db.Close() + defer store.Close() - feeder := NewFeedGenerator(db) + feeder := NewFeedGenerator(store) feedDidBase := os.Getenv("FEED_DID_BASE") if feedDidBase == "" { diff --git a/store/database.go b/store/database.go new file mode 100644 index 0000000..c5dd212 --- /dev/null +++ b/store/database.go @@ -0,0 +1,64 @@ +package store + +import ( + "database/sql" + "errors" + "fmt" + "log/slog" + "os" + + _ "github.com/glebarez/go-sqlite" +) + +type Store struct { + db *sql.DB +} + +func New(dbPath string) (*Store, error) { + err := createDbFile(dbPath) + if err != nil { + return nil, fmt.Errorf("create db file: %w", err) + } + + db, err := sql.Open("sqlite", dbPath) + if err != nil { + return nil, fmt.Errorf("open database: %w", err) + } + + err = db.Ping() + if err != nil { + return nil, fmt.Errorf("ping db: %w", err) + } + + err = createFeedTable(db) + if err != nil { + return nil, fmt.Errorf("creating feed table: %w", err) + } + + err = createSubscriptionsTable(db) + if err != nil { + return nil, fmt.Errorf("creating subscription table: %w", err) + } + + return &Store{db: db}, nil +} + +func (s *Store) Close() { + err := s.db.Close() + if err != nil { + slog.Error("failed to close db", "error", err) + } +} + +func createDbFile(dbFilename string) error { + if _, err := os.Stat(dbFilename); !errors.Is(err, os.ErrNotExist) { + return nil + } + + f, err := os.Create(dbFilename) + if err != nil { + return fmt.Errorf("create db file : %w", err) + } + f.Close() + return nil +} diff --git a/store/feed.go b/store/feed.go new file mode 100644 index 0000000..81deb50 --- /dev/null +++ b/store/feed.go @@ -0,0 +1,83 @@ +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, + 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 FeedItem struct { + ID int + ReplyURI string + UserDID string + SubscribedPostURI string +} + +func (s *Store) AddFeedItem(feedItem FeedItem) error { + sql := `INSERT INTO feed (replyURI, userDID, subscribedPostURI) VALUES (?, ?, ?) ON CONFLICT(replyURI, userDID) DO NOTHING;` + _, err := s.db.Exec(sql, feedItem.ReplyURI, feedItem.UserDID, feedItem.SubscribedPostURI) + if err != nil { + return fmt.Errorf("exec insert feed item: %w", err) + } + return nil +} + +func (s *Store) GetUsersFeedItems(usersDID string) ([]FeedItem, error) { + sql := "SELECT id, replyURI, userDID FROM feed WHERE userDID = ?;" + rows, err := s.db.Query(sql, usersDID) + if err != nil { + return nil, fmt.Errorf("run query to get users feed item: %w", err) + } + defer rows.Close() + + feedItems := make([]FeedItem, 0) + for rows.Next() { + var feedItem FeedItem + if err := rows.Scan(&feedItem.ID, &feedItem.ReplyURI, &feedItem.UserDID); err != nil { + return nil, fmt.Errorf("scan row: %w", err) + } + feedItems = append(feedItems, feedItem) + } + + return feedItems, nil +} + +func (s *Store) DeleteFeedItemsForSubscribedPostURIandUserDID(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 items: %w", err) + } + res, err := statement.Exec(subscribedPostURI, userDID) + if err != nil { + return fmt.Errorf("exec delete feed items: %w", err) + } + + n, _ := res.RowsAffected() + + slog.Info("delete feed res", "affected rows", n) + return nil +} diff --git a/store/subscription.go b/store/subscription.go new file mode 100644 index 0000000..5182be4 --- /dev/null +++ b/store/subscription.go @@ -0,0 +1,96 @@ +package store + +import ( + "database/sql" + "fmt" + "log/slog" +) + +func createSubscriptionsTable(db *sql.DB) error { + createSubscriptionsTableSQL := `CREATE TABLE IF NOT EXISTS subscriptions ( + "id" integer NOT NULL PRIMARY KEY AUTOINCREMENT, + "subscribedPostURI" TEXT, + "userDID" TEXT, + "subscriptionPostRkey" TEXT, + UNIQUE(subscribedPostURI, userDID) + );` + + slog.Info("Create subscriptions table...") + statement, err := db.Prepare(createSubscriptionsTableSQL) + if err != nil { + return fmt.Errorf("prepare DB statement to create subscriptions table: %w", err) + } + _, err = statement.Exec() + if err != nil { + return fmt.Errorf("exec sql statement to create subscriptions table: %w", err) + } + slog.Info("subscriptions table created") + + return nil +} + +type Subscription struct { + ID int + SubscribedPostURI string + UserDID string + SubscriptionPostRkey string +} + +func (s *Store) GetSubscriptionsForPost(postURI string) ([]string, error) { + sql := "SELECT userDID FROM subscriptions WHERE subscribedPostURI = ?" + rows, err := s.db.Query(sql, postURI) + if err != nil { + return nil, fmt.Errorf("run query to get subscriptions: %w", err) + } + defer rows.Close() + + dids := make([]string, 0) + for rows.Next() { + var subscription Subscription + if err := rows.Scan(&subscription.UserDID); err != nil { + return nil, fmt.Errorf("scan row: %w", err) + } + dids = append(dids, subscription.UserDID) + } + + return dids, nil +} + +func (s *Store) AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionRkey string) error { + sql := `INSERT INTO subscriptions (subscribedPostURI, userDID, subscriptionRkey) VALUES (?, ?, ?) ON CONFLICT(subscribedPostURI, userDID) DO NOTHING;` + _, err := s.db.Exec(sql, subscribedPostURI, userDid, subscriptionRkey) + if err != nil { + return fmt.Errorf("exec insert subscrptions: %w", err) + } + return nil +} + +func (s *Store) GetSubscribedPostURI(userDID, rkey string) (string, error) { + sql := "SELECT id, subscribedPostURI FROM subscriptions WHERE subscriptionRkey = ? AND userDID = ?;" + rows, err := s.db.Query(sql, rkey, userDID) + if err != nil { + return "", fmt.Errorf("run query to get subscribed post URI: %w", err) + } + defer rows.Close() + + subscribedPostURI := "" + for rows.Next() { + var subscription Subscription + if err := rows.Scan(&subscription.ID, &subscription.SubscribedPostURI); err != nil { + return "", fmt.Errorf("scan row: %w", err) + } + + subscribedPostURI = subscription.SubscribedPostURI + break + } + return subscribedPostURI, nil +} + +func (s *Store) DeleteSubscriptionForUser(userDID, postURI string) error { + sql := "DELETE FROM subscriptions WHERE subscribedPostURI = ? AND userDID = ?;" + _, err := s.db.Exec(sql, postURI, userDID) + if err != nil { + return fmt.Errorf("exec delete subscription for user: %w", err) + } + return nil +} -- 2.51.2 From 2a60b225e781703739e265c9eda1924db8e86fbc Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Mon, 18 Nov 2024 20:13:24 +0000 Subject: [PATCH 2/3] fix the subscriptionPostRkey rename --- store/subscription.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/store/subscription.go b/store/subscription.go index 5182be4..226faac 100644 --- a/store/subscription.go +++ b/store/subscription.go @@ -57,7 +57,7 @@ func (s *Store) GetSubscriptionsForPost(postURI string) ([]string, error) { } func (s *Store) AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionRkey string) error { - sql := `INSERT INTO subscriptions (subscribedPostURI, userDID, subscriptionRkey) VALUES (?, ?, ?) ON CONFLICT(subscribedPostURI, userDID) DO NOTHING;` + sql := `INSERT INTO subscriptions (subscribedPostURI, userDID, subscriptionPostRkey) VALUES (?, ?, ?) ON CONFLICT(subscribedPostURI, userDID) DO NOTHING;` _, err := s.db.Exec(sql, subscribedPostURI, userDid, subscriptionRkey) if err != nil { return fmt.Errorf("exec insert subscrptions: %w", err) @@ -66,7 +66,7 @@ func (s *Store) AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionR } func (s *Store) GetSubscribedPostURI(userDID, rkey string) (string, error) { - sql := "SELECT id, subscribedPostURI FROM subscriptions WHERE subscriptionRkey = ? AND userDID = ?;" + sql := "SELECT id, subscribedPostURI FROM subscriptions WHERE subscriptionPostRkey = ? AND userDID = ?;" rows, err := s.db.Query(sql, rkey, userDID) if err != nil { return "", fmt.Errorf("run query to get subscribed post URI: %w", err) -- 2.51.2 From 9697d99382131955a9f819afa04b6bbd2495975d Mon Sep 17 00:00:00 2001 From: Will Andrews Date: Mon, 18 Nov 2024 20:29:56 +0000 Subject: [PATCH 3/3] a few more renaming variables --- consumer.go | 4 ++-- store/subscription.go | 8 ++++---- 2 files changed, 6 insertions(+), 6 deletions(-) diff --git a/consumer.go b/consumer.go index 745fc23..271602a 100644 --- a/consumer.go +++ b/consumer.go @@ -154,8 +154,8 @@ func (h *handler) handleDeleteEvent(_ context.Context, event *models.Event) erro return nil } -func (h *handler) addDidToSubscribedPost(subscribedPostURI, userDid, rkey string) error { - err := h.store.AddSubscriptionForPost(subscribedPostURI, userDid, rkey) +func (h *handler) addDidToSubscribedPost(subscribedPostURI, userDid, subscriptionPostRkey string) error { + err := h.store.AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionPostRkey) if err != nil { return fmt.Errorf("add subscription for post: %w", err) } diff --git a/store/subscription.go b/store/subscription.go index 226faac..e56784a 100644 --- a/store/subscription.go +++ b/store/subscription.go @@ -56,18 +56,18 @@ func (s *Store) GetSubscriptionsForPost(postURI string) ([]string, error) { return dids, nil } -func (s *Store) AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionRkey string) error { +func (s *Store) AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionPostRkey string) error { sql := `INSERT INTO subscriptions (subscribedPostURI, userDID, subscriptionPostRkey) VALUES (?, ?, ?) ON CONFLICT(subscribedPostURI, userDID) DO NOTHING;` - _, err := s.db.Exec(sql, subscribedPostURI, userDid, subscriptionRkey) + _, err := s.db.Exec(sql, subscribedPostURI, userDid, subscriptionPostRkey) if err != nil { return fmt.Errorf("exec insert subscrptions: %w", err) } return nil } -func (s *Store) GetSubscribedPostURI(userDID, rkey string) (string, error) { +func (s *Store) GetSubscribedPostURI(userDID, subscriptionPostRkey string) (string, error) { sql := "SELECT id, subscribedPostURI FROM subscriptions WHERE subscriptionPostRkey = ? AND userDID = ?;" - rows, err := s.db.Query(sql, rkey, userDID) + rows, err := s.db.Query(sql, subscriptionPostRkey, userDID) if err != nil { return "", fmt.Errorf("run query to get subscribed post URI: %w", err) }