diff --git a/appview/bsky/client.go b/appview/bsky/client.go new file mode 100644 --- /dev/null +++ b/appview/bsky/client.go @@ -0,0 +1,43 @@ +package bsky + +import ( + "context" + "log" + + bsky "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/xrpc" + "tangled.org/core/appview/models" + "tangled.org/core/consts" +) + +func FetchPosts(ctx context.Context, c *xrpc.Client, limit int, cursor string) ([]models.BskyPost, string, error) { + resp, err := bsky.FeedGetAuthorFeed(ctx, c, consts.TangledDid, cursor, "posts_no_replies", false, int64(limit)) + if err != nil { + return nil, "", err + } + + var posts []models.BskyPost + for _, feedViewPost := range resp.Feed { + // skip quote posts + if feedViewPost.Post.Embed != nil && feedViewPost.Post.Embed.EmbedRecord_View != nil { + continue + } + + post, err := models.NewBskyPostFromView(feedViewPost.Post) + if err != nil { + log.Println(err) + continue + } + + posts = append(posts, *post) + } + + return posts, stringPtr(resp.Cursor), nil +} + +func stringPtr(s *string) string { + if s == nil { + return "" + } + return *s +} diff --git a/appview/config/config.go b/appview/config/config.go --- a/appview/config/config.go +++ b/appview/config/config.go @@ -100,6 +100,10 @@ GoodFirstIssue string `env:"GFI, default=at://did:plc:wshs7t2adsemcrrd4snkeqli/sh.tangled.label.definition/good-first-issue"` } +type BlueskyConfig struct { + UpdateInterval time.Duration `env:"UPDATE_INTERVAL, default=1h"` +} + func (cfg RedisConfig) ToURL() string { u := &url.URL{ Scheme: "redis", @@ -129,6 +133,7 @@ Pds PdsConfig `env:",prefix=TANGLED_PDS_"` Cloudflare Cloudflare `env:",prefix=TANGLED_CLOUDFLARE_"` Label LabelConfig `env:",prefix=TANGLED_LABEL_"` + Bluesky BlueskyConfig `env:",prefix=TANGLED_BLUESKY_"` } func LoadConfig(ctx context.Context) (*Config, error) { diff --git a/appview/db/bsky.go b/appview/db/bsky.go new file mode 100644 --- /dev/null +++ b/appview/db/bsky.go @@ -0,0 +1,116 @@ +package db + +import ( + "database/sql" + "encoding/json" + "time" + + "tangled.org/core/appview/models" +) + +func InsertBlueskyPosts(e Execer, posts []models.BskyPost) error { + if len(posts) == 0 { + return nil + } + + stmt, err := e.Prepare(` + insert or replace into bluesky_posts (rkey, text, created_at, langs, facets, embed, like_count, reply_count, repost_count, quote_count) + values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + `) + if err != nil { + return err + } + defer stmt.Close() + + for _, post := range posts { + var langsJSON, facetsJSON, embedJSON []byte + + if len(post.Langs) > 0 { + langsJSON, _ = json.Marshal(post.Langs) + } + if len(post.Facets) > 0 { + facetsJSON, _ = json.Marshal(post.Facets) + } + if post.Embed != nil { + embedJSON, _ = json.Marshal(post.Embed) + } + + _, err := stmt.Exec( + post.Rkey, + post.Text, + post.CreatedAt.Format(time.RFC3339), + nullString(langsJSON), + nullString(facetsJSON), + nullString(embedJSON), + post.LikeCount, + post.ReplyCount, + post.RepostCount, + post.QuoteCount, + ) + if err != nil { + return err + } + } + + return nil +} + +func nullString(b []byte) any { + if len(b) == 0 { + return nil + } + return string(b) +} + +func GetBlueskyPosts(e Execer, limit int) ([]models.BskyPost, error) { + query := ` + select rkey, text, created_at, langs, facets, embed, like_count, reply_count, repost_count, quote_count + from bluesky_posts + order by created_at desc + limit ? + ` + + rows, err := e.Query(query, limit) + if err != nil { + return nil, err + } + defer rows.Close() + + var posts []models.BskyPost + for rows.Next() { + var rkey, text, createdAt string + var langs, facets, embed sql.Null[string] + var likeCount, replyCount, repostCount, quoteCount int64 + + err := rows.Scan(&rkey, &text, &createdAt, &langs, &facets, &embed, &likeCount, &replyCount, &repostCount, "eCount) + if err != nil { + return nil, err + } + + post := models.BskyPost{ + Rkey: rkey, + Text: text, + LikeCount: likeCount, + ReplyCount: replyCount, + RepostCount: repostCount, + QuoteCount: quoteCount, + } + + if t, err := time.Parse(time.RFC3339, createdAt); err == nil { + post.CreatedAt = t + } + if langs.Valid && langs.V != "" { + json.Unmarshal([]byte(langs.V), &post.Langs) + } + if facets.Valid && facets.V != "" { + json.Unmarshal([]byte(facets.V), &post.Facets) + } + if embed.Valid && embed.V != "" { + json.Unmarshal([]byte(embed.V), &post.Embed) + } + + posts = append(posts, post) + } + + return posts, rows.Err() +} diff --git a/appview/db/db.go b/appview/db/db.go --- a/appview/db/db.go +++ b/appview/db/db.go @@ -596,6 +596,19 @@ foreign key (webhook_id) references webhooks(id) on delete cascade ); + create table if not exists bluesky_posts ( + rkey text primary key, + text text not null, + created_at text not null, + langs text, + facets text, + embed text, + like_count integer not null default 0, + reply_count integer not null default 0, + repost_count integer not null default 0, + quote_count integer not null default 0 + ); + create table if not exists migrations ( id integer primary key autoincrement, name text unique diff --git a/appview/models/bsky.go b/appview/models/bsky.go new file mode 100644 --- /dev/null +++ b/appview/models/bsky.go @@ -0,0 +1,72 @@ +package models + +import ( + "time" + + apibsky "github.com/bluesky-social/indigo/api/bsky" + "github.com/bluesky-social/indigo/atproto/syntax" +) + +type BskyPost struct { + Rkey string + Text string + CreatedAt time.Time + Langs []string + Tags []string + Embed *apibsky.FeedDefs_PostView_Embed + Facets []*apibsky.RichtextFacet + Labels *apibsky.FeedPost_Labels + Reply *apibsky.FeedPost_ReplyRef + LikeCount int64 + ReplyCount int64 + RepostCount int64 + QuoteCount int64 +} + +func NewBskyPostFromView(postView *apibsky.FeedDefs_PostView) (*BskyPost, error) { + atUri, err := syntax.ParseATURI(postView.Uri) + if err != nil { + return nil, err + } + + // decode the record to get FeedPost + feedPost, ok := postView.Record.Val.(*apibsky.FeedPost) + if !ok { + return nil, err + } + + createdAt, err := time.Parse(time.RFC3339, feedPost.CreatedAt) + if err != nil { + return nil, err + } + + var likeCount, replyCount, repostCount, quoteCount int64 + if postView.LikeCount != nil { + likeCount = *postView.LikeCount + } + if postView.ReplyCount != nil { + replyCount = *postView.ReplyCount + } + if postView.RepostCount != nil { + repostCount = *postView.RepostCount + } + if postView.QuoteCount != nil { + quoteCount = *postView.QuoteCount + } + + return &BskyPost{ + Rkey: atUri.RecordKey().String(), + Text: feedPost.Text, + CreatedAt: createdAt, + Langs: feedPost.Langs, + Tags: feedPost.Tags, + Embed: postView.Embed, + Facets: feedPost.Facets, + Labels: feedPost.Labels, + Reply: feedPost.Reply, + LikeCount: likeCount, + ReplyCount: replyCount, + RepostCount: repostCount, + QuoteCount: quoteCount, + }, nil +} diff --git a/appview/oauth/handler.go b/appview/oauth/handler.go --- a/appview/oauth/handler.go +++ b/appview/oauth/handler.go @@ -19,6 +19,7 @@ "tangled.org/core/appview/db" "tangled.org/core/appview/models" "tangled.org/core/consts" + "tangled.org/core/idresolver" "tangled.org/core/orm" "tangled.org/core/tid" ) @@ -129,7 +130,7 @@ } l.Debug("adding to default spindle") - session, err := o.createAppPasswordSession(o.Config.Core.AppPassword, consts.TangledDid) + session, err := CreateAppPasswordSession(o.IdResolver, o.Config.Core.AppPassword, consts.TangledDid) if err != nil { l.Error("failed to create session", "err", err) return @@ -168,7 +169,7 @@ } l.Debug("adding to default knot") - session, err := o.createAppPasswordSession(o.Config.Core.TmpAltAppPassword, consts.IcyDid) + session, err := CreateAppPasswordSession(o.IdResolver, o.Config.Core.TmpAltAppPassword, consts.IcyDid) if err != nil { l.Error("failed to create session", "err", err) return @@ -241,19 +242,19 @@ l.Debug("successfully created empty Tangled profile on PDS and DB") } -// create a session using apppasswords -type session struct { +// create a AppPasswordSession using apppasswords +type AppPasswordSession struct { AccessJwt string `json:"accessJwt"` PdsEndpoint string Did string } -func (o *OAuth) createAppPasswordSession(appPassword, did string) (*session, error) { +func CreateAppPasswordSession(res *idresolver.Resolver, appPassword, did string) (*AppPasswordSession, error) { if appPassword == "" { - return nil, fmt.Errorf("no app password configured, skipping member addition") + return nil, fmt.Errorf("no app password configured") } - resolved, err := o.IdResolver.ResolveIdent(context.Background(), did) + resolved, err := res.ResolveIdent(context.Background(), did) if err != nil { return nil, fmt.Errorf("failed to resolve tangled.sh DID %s: %v", did, err) } @@ -290,7 +291,7 @@ return nil, fmt.Errorf("failed to create session: HTTP %d", sessionResp.StatusCode) } - var session session + var session AppPasswordSession if err := json.NewDecoder(sessionResp.Body).Decode(&session); err != nil { return nil, fmt.Errorf("failed to decode session response: %v", err) } @@ -301,7 +302,7 @@ return &session, nil } -func (s *session) putRecord(record any, collection string) error { +func (s *AppPasswordSession) putRecord(record any, collection string) error { recordBytes, err := json.Marshal(record) if err != nil { return fmt.Errorf("failed to marshal knot member record: %w", err) diff --git a/appview/pages/pages.go b/appview/pages/pages.go --- a/appview/pages/pages.go +++ b/appview/pages/pages.go @@ -339,6 +339,7 @@ Timeline []models.TimelineEvent Repos []models.Repo GfiLabel *models.LabelDefinition + BlueskyPosts []models.BskyPost } func (p *Pages) Timeline(w io.Writer, params TimelineParams) error { diff --git a/appview/state/state.go b/appview/state/state.go --- a/appview/state/state.go +++ b/appview/state/state.go @@ -12,6 +12,7 @@ "tangled.org/core/api/tangled" "tangled.org/core/appview" + "tangled.org/core/appview/bsky" "tangled.org/core/appview/config" "tangled.org/core/appview/db" "tangled.org/core/appview/indexer" @@ -25,6 +26,7 @@ "tangled.org/core/appview/reporesolver" "tangled.org/core/appview/validator" xrpcclient "tangled.org/core/appview/xrpcclient" + "tangled.org/core/consts" "tangled.org/core/eventconsumer" "tangled.org/core/idresolver" "tangled.org/core/jetstream" @@ -38,6 +40,7 @@ atpclient "github.com/bluesky-social/indigo/atproto/client" "github.com/bluesky-social/indigo/atproto/syntax" lexutil "github.com/bluesky-social/indigo/lex/util" + "github.com/bluesky-social/indigo/xrpc" securejoin "github.com/cyphar/filepath-securejoin" "github.com/go-chi/chi/v5" "github.com/posthog/posthog-go" @@ -198,6 +201,9 @@ logger, validator, } + + // fetch initial bluesky posts if configured + go fetchBskyPosts(ctx, res, config, d, logger) return state, nil }