diff --git a/consumer.go b/consumer.go index 25aeb31..98142bb 100644 --- a/consumer.go +++ b/consumer.go @@ -3,33 +3,21 @@ package main import ( "context" - "encoding/json" "fmt" "log/slog" - "strings" "time" - apibsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/jetstream/pkg/client" "github.com/bluesky-social/jetstream/pkg/client/schedulers/sequential" - "github.com/bluesky-social/jetstream/pkg/models" - "github.com/bugsnag/bugsnag-go/v2" ) -type ConsumerStore interface { - GetSubscriptionsForPost(postURI string) ([]string, error) - AddSubscriptionForPost(subscribedPostURI, userDid, subscriptionPostRkey string) error - GetSubscribedPostURI(userDID, subscriptionPostRkey string) (string, error) - DeleteSubscriptionForUser(userDID, postURI string) error - DeleteFeedPostsForSubscribedPostURIandUserDID(subscribedPostURI, userDID string) error -} - type consumer struct { - cfg *client.ClientConfig - store ConsumerStore + cfg *client.ClientConfig + handler *handler + logger *slog.Logger } -func NewConsumer(jsAddr string, store ConsumerStore) *consumer { +func NewConsumer(jsAddr string, logger *slog.Logger, handler *handler) *consumer { cfg := client.DefaultClientConfig() if jsAddr != "" { cfg.WebsocketURL = jsAddr @@ -38,145 +26,29 @@ func NewConsumer(jsAddr string, store ConsumerStore) *consumer { "app.bsky.feed.post", } cfg.WantedDids = []string{} + return &consumer{ - cfg: cfg, - store: store, + cfg: cfg, + logger: logger, + handler: handler, } } -func (con *consumer) Consume(ctx context.Context, feedGen *FeedGenerator, logger *slog.Logger) error { - h := &handler{ - feedGenerator: feedGen, - store: con.store, - } - - scheduler := sequential.NewScheduler("jetstream_localdev", logger, h.HandleEvent) +func (c *consumer) Consume(ctx context.Context) error { + scheduler := sequential.NewScheduler("jetstream_localdev", c.logger, c.handler.HandleEvent) defer scheduler.Shutdown() - c, err := client.NewClient(con.cfg, logger, scheduler) + client, err := client.NewClient(c.cfg, c.logger, scheduler) if err != nil { return fmt.Errorf("failed to create client: %w", err) } cursor := time.Now().Add(1 * -time.Minute).UnixMicro() - if err := c.ConnectAndRead(ctx, &cursor); err != nil { + if err := client.ConnectAndRead(ctx, &cursor); err != nil { return fmt.Errorf("connect and read: %w", err) } slog.Info("stopping consume") return nil } - -type handler struct { - feedGenerator *FeedGenerator - store ConsumerStore -} - -func (h *handler) HandleEvent(ctx context.Context, event *models.Event) error { - if event.Commit == nil { - return nil - } - - switch event.Commit.Operation { - case models.CommitOperationCreate: - return h.handleCreateEvent(ctx, event) - case models.CommitOperationDelete: - return h.handleDeleteEvent(ctx, event) - default: - return nil - } -} - -func (h *handler) handleCreateEvent(_ context.Context, event *models.Event) error { - if event.Commit.Collection != "app.bsky.feed.post" { - return nil - } - - var post apibsky.FeedPost - if err := json.Unmarshal(event.Commit.Record, &post); err != nil { - // ignore this - return nil - } - - // we only care about posts that have parents which are replies - if post.Reply == nil || post.Reply.Parent == nil || post.Reply.Parent.Uri == "" { - return nil - } - - subscribedPostURI := post.Reply.Parent.Uri - - // 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") { - // For now just look for me - if event.Did != "did:plc:dadhhalkfcq3gucaq25hjqon" { - return nil - } - 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.getSubscribedDidsForPost(subscribedPostURI) - if len(subscribedDids) == 0 { - return nil - } - - slog.Info("post is a reply to a post that users are subscribed to", "subscribed post URI", subscribedPostURI, "dids", subscribedDids, "RKey", 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 -} - -func (h *handler) handleDeleteEvent(_ context.Context, event *models.Event) error { - if event.Commit.Collection != "app.bsky.feed.post" { - return nil - } - - // temp ignore everyone but me - if event.Did != "did:plc:dadhhalkfcq3gucaq25hjqon" { - return nil - } - slog.Info("delete event received", "did", event.Did, "rkey", event.Commit.RKey) - subscribedPostURI, err := h.store.GetSubscribedPostURI(event.Did, event.Commit.RKey) - if err != nil { - 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) - } - - // 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 = h.store.DeleteFeedPostsForSubscribedPostURIandUserDID(subscribedPostURI, event.Did) - if err != nil { - 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 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, "subscribedPostURI", subscribedPostURI, "user DID", event.Did) - return fmt.Errorf("delete subscription and user: %w", err) - } - - return nil -} - -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) - } - return nil -} - -func (h *handler) getSubscribedDidsForPost(postURI string) []string { - dids, err := h.store.GetSubscriptionsForPost(postURI) - if err != nil { - slog.Error("getting subscriptions for post", "error", err) - bugsnag.Notify(err) - } - - return dids -} diff --git a/feedgenerator.go b/feedgenerator.go index e19f95e..fb27641 100644 --- a/feedgenerator.go +++ b/feedgenerator.go @@ -3,14 +3,11 @@ package main import ( "context" "fmt" - "log/slog" - "github.com/bugsnag/bugsnag-go/v2" "github.com/willdot/bskyfeedgen/store" ) type feedStore interface { - AddFeedPost(feedItem store.FeedPost) error GetUsersFeed(usersDID string) ([]store.FeedPost, error) } @@ -46,19 +43,3 @@ func (f *FeedGenerator) GetFeed(ctx context.Context, userDID, feed, cursor strin return resp, nil } - -func (f *FeedGenerator) AddToFeedPosts(usersDids []string, subscribedPostURI, replyPostURI string) { - for _, did := range usersDids { - feedItem := store.FeedPost{ - ReplyURI: replyPostURI, - UserDID: did, - SubscribedPostURI: subscribedPostURI, - } - err := f.store.AddFeedPost(feedItem) - if err != nil { - 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 a8357c3..bd725b9 100644 --- a/main.go +++ b/main.go @@ -71,7 +71,7 @@ func main() { if enableJS == "true" { slog.Info("enabling jetstream consume") - go consumeLoop(ctx, jsServerAddr, feeder, store) + go consumeLoop(ctx, jsServerAddr, store) } server := NewServer(443, feeder, feedHost, feedDidBase) @@ -87,11 +87,14 @@ func main() { time.Sleep(time.Second) } -func consumeLoop(ctx context.Context, jsServerAddr string, feeder *FeedGenerator, store *store.Store) { - consumer := NewConsumer(jsServerAddr, store) +func consumeLoop(ctx context.Context, jsServerAddr string, store *store.Store) { + handler := handler{ + store: store, + } + consumer := NewConsumer(jsServerAddr, slog.Default(), &handler) retry.Do(func() error { - err := consumer.Consume(ctx, feeder, slog.Default()) + err := consumer.Consume(ctx) if err != nil { if errors.Is(err, context.Canceled) { return nil