nandi/gleanpublic⑂ Fork 0
⑂ 2c8faf0
Commits
⬇ Clone ▾
git clone https://git.rickub.com/nandi/glean.git
git clone ssh://git@rickub.com/nandi/glean.git

Host key fingerprint (ed25519): SHA256:iycHnxEyq0Q7uyVpB7JlznP0G7JrTPXLYRcAU5CSLhc — verify it before your first connect.

add background sync and record synchronization supportUnverified

Julien Robert committed 2026-04-20T15:11:43+02:00 Browse files
2c8faf0 parent: 7c98e02
modified .gitignore +3 -0
@@ -28,3 +28,6 @@ go.work.sum
2828 # tailwind
2929 static/output.css
3030 node_modules/
31+
32+# todo
33+todo.md
@@ -28,3 +28,6 @@ go.work.sum
28 # tailwind28 # tailwind
29 static/output.css29 static/output.css
30 node_modules/30 node_modules/
31+
32+# todo
33+todo.md
modified internal/atproto/client.go +56 -4
@@ -88,9 +88,9 @@ func (c *Client) createRecordWithAPI(ctx context.Context, did, collection string
8888 CID string `json:"cid"`
8989 }
9090
91- nsid, err := syntax.ParseNSID(collection)
91+ nsid, err := syntax.ParseNSID("com.atproto.repo.createRecord")
9292 if err != nil {
93- return "", "", fmt.Errorf("parsing collection NSID: %w", err)
93+ return "", "", fmt.Errorf("parsing NSID: %w", err)
9494 }
9595
9696 if err := c.APIClient.Post(ctx, nsid, input, &out); err != nil {
@@ -100,13 +100,21 @@ func (c *Client) createRecordWithAPI(ctx context.Context, did, collection string
100100 }
101101
102102 func (c *Client) DeleteRecord(ctx context.Context, did, collection, rkey string) error {
103- body := map[string]any{
103+ input := map[string]any{
104104 "repo": did,
105105 "collection": collection,
106106 "rkey": rkey,
107107 }
108108
109- data, err := json.Marshal(body)
109+ if c.APIClient != nil {
110+ nsid, err := syntax.ParseNSID("com.atproto.repo.deleteRecord")
111+ if err != nil {
112+ return fmt.Errorf("parsing NSID: %w", err)
113+ }
114+ return c.APIClient.Post(ctx, nsid, input, nil)
115+ }
116+
117+ data, err := json.Marshal(input)
110118 if err != nil {
111119 return err
112120 }
@@ -133,6 +141,10 @@ func (c *Client) DeleteRecord(ctx context.Context, did, collection, rkey string)
133141 }
134142
135143 func (c *Client) ListRecords(ctx context.Context, did, collection string, limit int, cursor string) ([]Record, string, error) {
144+ if c.APIClient != nil {
145+ return c.listRecordsWithAPI(ctx, did, collection, limit, cursor)
146+ }
147+
136148 url := fmt.Sprintf("%s/xrpc/com.atproto.repo.listRecords?repo=%s&collection=%s", c.pdsURL, did, collection)
137149 if limit > 0 {
138150 url += fmt.Sprintf("&limit=%d", limit)
@@ -177,6 +189,46 @@ func (c *Client) ListRecords(ctx context.Context, did, collection string, limit
177189 Value: r.Value,
178190 }
179191 }
192+ return records, result.Cursor, nil
193+}
194+
195+func (c *Client) listRecordsWithAPI(ctx context.Context, did, collection string, limit int, cursor string) ([]Record, string, error) {
196+ nsid, err := syntax.ParseNSID("com.atproto.repo.listRecords")
197+ if err != nil {
198+ return nil, "", fmt.Errorf("parsing NSID: %w", err)
199+ }
180200
201+ params := map[string]any{
202+ "repo": did,
203+ "collection": collection,
204+ }
205+ if limit > 0 {
206+ params["limit"] = limit
207+ }
208+ if cursor != "" {
209+ params["cursor"] = cursor
210+ }
211+
212+ var result struct {
213+ Records []struct {
214+ URI string `json:"uri"`
215+ CID string `json:"cid"`
216+ Value json.RawMessage `json:"value"`
217+ } `json:"records"`
218+ Cursor string `json:"cursor"`
219+ }
220+
221+ if err := c.APIClient.Get(ctx, nsid, params, &result); err != nil {
222+ return nil, "", err
223+ }
224+
225+ records := make([]Record, len(result.Records))
226+ for i, r := range result.Records {
227+ records[i] = Record{
228+ URI: r.URI,
229+ CID: r.CID,
230+ Value: r.Value,
231+ }
232+ }
181233 return records, result.Cursor, nil
182234 }
@@ -88,9 +88,9 @@ func (c *Client) createRecordWithAPI(ctx context.Context, did, collection string
88 CID string `json:"cid"`88 CID string `json:"cid"`
89 }89 }
90 90
91- nsid, err := syntax.ParseNSID(collection)91+ nsid, err := syntax.ParseNSID("com.atproto.repo.createRecord")
92 if err != nil {92 if err != nil {
93- return "", "", fmt.Errorf("parsing collection NSID: %w", err)93+ return "", "", fmt.Errorf("parsing NSID: %w", err)
94 }94 }
95 95
96 if err := c.APIClient.Post(ctx, nsid, input, &out); err != nil {96 if err := c.APIClient.Post(ctx, nsid, input, &out); err != nil {
@@ -100,13 +100,21 @@ func (c *Client) createRecordWithAPI(ctx context.Context, did, collection string
100 }100 }
101 101
102 func (c *Client) DeleteRecord(ctx context.Context, did, collection, rkey string) error {102 func (c *Client) DeleteRecord(ctx context.Context, did, collection, rkey string) error {
103- body := map[string]any{103+ input := map[string]any{
104 "repo": did,104 "repo": did,
105 "collection": collection,105 "collection": collection,
106 "rkey": rkey,106 "rkey": rkey,
107 }107 }
108 108
109- data, err := json.Marshal(body)109+ if c.APIClient != nil {
110+ nsid, err := syntax.ParseNSID("com.atproto.repo.deleteRecord")
111+ if err != nil {
112+ return fmt.Errorf("parsing NSID: %w", err)
113+ }
114+ return c.APIClient.Post(ctx, nsid, input, nil)
115+ }
116+
117+ data, err := json.Marshal(input)
110 if err != nil {118 if err != nil {
111 return err119 return err
112 }120 }
@@ -133,6 +141,10 @@ func (c *Client) DeleteRecord(ctx context.Context, did, collection, rkey string)
133 }141 }
134 142
135 func (c *Client) ListRecords(ctx context.Context, did, collection string, limit int, cursor string) ([]Record, string, error) {143 func (c *Client) ListRecords(ctx context.Context, did, collection string, limit int, cursor string) ([]Record, string, error) {
144+ if c.APIClient != nil {
145+ return c.listRecordsWithAPI(ctx, did, collection, limit, cursor)
146+ }
147+
136 url := fmt.Sprintf("%s/xrpc/com.atproto.repo.listRecords?repo=%s&collection=%s", c.pdsURL, did, collection)148 url := fmt.Sprintf("%s/xrpc/com.atproto.repo.listRecords?repo=%s&collection=%s", c.pdsURL, did, collection)
137 if limit > 0 {149 if limit > 0 {
138 url += fmt.Sprintf("&limit=%d", limit)150 url += fmt.Sprintf("&limit=%d", limit)
@@ -177,6 +189,46 @@ func (c *Client) ListRecords(ctx context.Context, did, collection string, limit
177 Value: r.Value,189 Value: r.Value,
178 }190 }
179 }191 }
192+ return records, result.Cursor, nil
193+}
194+
195+func (c *Client) listRecordsWithAPI(ctx context.Context, did, collection string, limit int, cursor string) ([]Record, string, error) {
196+ nsid, err := syntax.ParseNSID("com.atproto.repo.listRecords")
197+ if err != nil {
198+ return nil, "", fmt.Errorf("parsing NSID: %w", err)
199+ }
180 200
201+ params := map[string]any{
202+ "repo": did,
203+ "collection": collection,
204+ }
205+ if limit > 0 {
206+ params["limit"] = limit
207+ }
208+ if cursor != "" {
209+ params["cursor"] = cursor
210+ }
211+
212+ var result struct {
213+ Records []struct {
214+ URI string `json:"uri"`
215+ CID string `json:"cid"`
216+ Value json.RawMessage `json:"value"`
217+ } `json:"records"`
218+ Cursor string `json:"cursor"`
219+ }
220+
221+ if err := c.APIClient.Get(ctx, nsid, params, &result); err != nil {
222+ return nil, "", err
223+ }
224+
225+ records := make([]Record, len(result.Records))
226+ for i, r := range result.Records {
227+ records[i] = Record{
228+ URI: r.URI,
229+ CID: r.CID,
230+ Value: r.Value,
231+ }
232+ }
181 return records, result.Cursor, nil233 return records, result.Cursor, nil
182 }234 }
added internal/atproto/firehose_handler.go +137 -0
new file mode 100644
@@ -0,0 +1,137 @@
1+package atproto
2+
3+import (
4+ "context"
5+ "encoding/json"
6+ "database/sql"
7+ "log/slog"
8+ "time"
9+
10+ "pkg.rbrt.fr/glean/internal/db"
11+)
12+
13+type FirehoseDBHandler struct {
14+ db *db.DB
15+ logger *slog.Logger
16+}
17+
18+func NewFirehoseDBHandler(database *db.DB, logger *slog.Logger) *FirehoseDBHandler {
19+ return &FirehoseDBHandler{db: database, logger: logger}
20+}
21+
22+func (h *FirehoseDBHandler) Handle(ctx context.Context, event *FirehoseEvent) error {
23+ switch event.Collection {
24+ case "at.glean.subscription":
25+ return h.handleSubscription(ctx, event)
26+ case "at.glean.like":
27+ return h.handleLike(ctx, event)
28+ case "at.glean.annotation":
29+ return h.handleAnnotation(ctx, event)
30+ }
31+ return nil
32+}
33+
34+func (h *FirehoseDBHandler) handleSubscription(ctx context.Context, event *FirehoseEvent) error {
35+ switch event.Type {
36+ case "create", "update":
37+ var rec SubscriptionRecord
38+ if err := json.Unmarshal(event.Value, &rec); err != nil {
39+ return err
40+ }
41+ if rec.FeedURL == "" {
42+ return nil
43+ }
44+
45+ existing, err := h.db.GetSubscription(ctx, event.DID, rec.FeedURL)
46+ if err == nil && existing != nil {
47+ if !existing.URI.Valid || existing.URI.String == "" {
48+ return h.db.UpdateSubscriptionURI(ctx, event.DID, rec.FeedURL, event.URI, event.CID)
49+ }
50+ return nil
51+ }
52+
53+ f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)}
54+ _ = h.db.UpsertFeed(ctx, f)
55+ return h.db.CreateSubscription(ctx, event.DID, rec.FeedURL, rec.Category, event.URI, event.CID)
56+
57+ case "delete":
58+ parsed, ok := ParseRecordURI(event.URI)
59+ if !ok {
60+ return nil
61+ }
62+ subs, _ := h.db.ListSubscriptions(ctx, event.DID, "", 100, 0)
63+ for _, sub := range subs {
64+ if sub.URI.Valid && sub.URI.String == event.URI {
65+ return h.db.DeleteSubscription(ctx, event.DID, sub.FeedURL)
66+ }
67+ }
68+ _ = parsed
69+ }
70+ return nil
71+}
72+
73+func (h *FirehoseDBHandler) handleLike(ctx context.Context, event *FirehoseEvent) error {
74+ switch event.Type {
75+ case "create":
76+ var rec LikeRecord
77+ if err := json.Unmarshal(event.Value, &rec); err != nil {
78+ return err
79+ }
80+ if rec.FeedURL == "" || rec.ArticleURL == "" {
81+ return nil
82+ }
83+
84+ exists, err := h.db.HasLiked(ctx, event.DID, rec.FeedURL, rec.ArticleURL)
85+ if err != nil || exists {
86+ return nil
87+ }
88+
89+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
90+ return h.db.CreateLike(ctx, &db.Like{
91+ URI: event.URI,
92+ AuthorDID: event.DID,
93+ FeedURL: rec.FeedURL,
94+ ArticleURL: rec.ArticleURL,
95+ CreatedAt: sql.NullTime{Time: t, Valid: true},
96+ CID: sql.NullString{String: event.CID, Valid: event.CID != ""},
97+ })
98+
99+ case "delete":
100+ return h.db.DeleteLike(ctx, event.URI)
101+ }
102+ return nil
103+}
104+
105+func (h *FirehoseDBHandler) handleAnnotation(ctx context.Context, event *FirehoseEvent) error {
106+ switch event.Type {
107+ case "create":
108+ var rec AnnotationRecord
109+ if err := json.Unmarshal(event.Value, &rec); err != nil {
110+ return err
111+ }
112+ if rec.FeedURL == "" || rec.ArticleURL == "" {
113+ return nil
114+ }
115+
116+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
117+ a := &db.Annotation{
118+ URI: event.URI,
119+ AuthorDID: event.DID,
120+ FeedURL: rec.FeedURL,
121+ ArticleURL: rec.ArticleURL,
122+ Quote: db.NullStr(rec.Quote),
123+ Note: db.NullStr(rec.Note),
124+ Tags: db.NullStrTags(rec.Tags),
125+ CreatedAt: sql.NullTime{Time: t, Valid: true},
126+ CID: sql.NullString{String: event.CID, Valid: event.CID != ""},
127+ }
128+ if rec.Rating > 0 {
129+ a.Rating = sql.NullInt64{Int64: int64(rec.Rating), Valid: true}
130+ }
131+ return h.db.CreateAnnotation(ctx, a)
132+
133+ case "delete":
134+ return h.db.DeleteAnnotation(ctx, event.URI)
135+ }
136+ return nil
137+}
new file mode 100644
@@ -0,0 +1,137 @@
1+package atproto
2+
3+import (
4+ "context"
5+ "encoding/json"
6+ "database/sql"
7+ "log/slog"
8+ "time"
9+
10+ "pkg.rbrt.fr/glean/internal/db"
11+)
12+
13+type FirehoseDBHandler struct {
14+ db *db.DB
15+ logger *slog.Logger
16+}
17+
18+func NewFirehoseDBHandler(database *db.DB, logger *slog.Logger) *FirehoseDBHandler {
19+ return &FirehoseDBHandler{db: database, logger: logger}
20+}
21+
22+func (h *FirehoseDBHandler) Handle(ctx context.Context, event *FirehoseEvent) error {
23+ switch event.Collection {
24+ case "at.glean.subscription":
25+ return h.handleSubscription(ctx, event)
26+ case "at.glean.like":
27+ return h.handleLike(ctx, event)
28+ case "at.glean.annotation":
29+ return h.handleAnnotation(ctx, event)
30+ }
31+ return nil
32+}
33+
34+func (h *FirehoseDBHandler) handleSubscription(ctx context.Context, event *FirehoseEvent) error {
35+ switch event.Type {
36+ case "create", "update":
37+ var rec SubscriptionRecord
38+ if err := json.Unmarshal(event.Value, &rec); err != nil {
39+ return err
40+ }
41+ if rec.FeedURL == "" {
42+ return nil
43+ }
44+
45+ existing, err := h.db.GetSubscription(ctx, event.DID, rec.FeedURL)
46+ if err == nil && existing != nil {
47+ if !existing.URI.Valid || existing.URI.String == "" {
48+ return h.db.UpdateSubscriptionURI(ctx, event.DID, rec.FeedURL, event.URI, event.CID)
49+ }
50+ return nil
51+ }
52+
53+ f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)}
54+ _ = h.db.UpsertFeed(ctx, f)
55+ return h.db.CreateSubscription(ctx, event.DID, rec.FeedURL, rec.Category, event.URI, event.CID)
56+
57+ case "delete":
58+ parsed, ok := ParseRecordURI(event.URI)
59+ if !ok {
60+ return nil
61+ }
62+ subs, _ := h.db.ListSubscriptions(ctx, event.DID, "", 100, 0)
63+ for _, sub := range subs {
64+ if sub.URI.Valid && sub.URI.String == event.URI {
65+ return h.db.DeleteSubscription(ctx, event.DID, sub.FeedURL)
66+ }
67+ }
68+ _ = parsed
69+ }
70+ return nil
71+}
72+
73+func (h *FirehoseDBHandler) handleLike(ctx context.Context, event *FirehoseEvent) error {
74+ switch event.Type {
75+ case "create":
76+ var rec LikeRecord
77+ if err := json.Unmarshal(event.Value, &rec); err != nil {
78+ return err
79+ }
80+ if rec.FeedURL == "" || rec.ArticleURL == "" {
81+ return nil
82+ }
83+
84+ exists, err := h.db.HasLiked(ctx, event.DID, rec.FeedURL, rec.ArticleURL)
85+ if err != nil || exists {
86+ return nil
87+ }
88+
89+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
90+ return h.db.CreateLike(ctx, &db.Like{
91+ URI: event.URI,
92+ AuthorDID: event.DID,
93+ FeedURL: rec.FeedURL,
94+ ArticleURL: rec.ArticleURL,
95+ CreatedAt: sql.NullTime{Time: t, Valid: true},
96+ CID: sql.NullString{String: event.CID, Valid: event.CID != ""},
97+ })
98+
99+ case "delete":
100+ return h.db.DeleteLike(ctx, event.URI)
101+ }
102+ return nil
103+}
104+
105+func (h *FirehoseDBHandler) handleAnnotation(ctx context.Context, event *FirehoseEvent) error {
106+ switch event.Type {
107+ case "create":
108+ var rec AnnotationRecord
109+ if err := json.Unmarshal(event.Value, &rec); err != nil {
110+ return err
111+ }
112+ if rec.FeedURL == "" || rec.ArticleURL == "" {
113+ return nil
114+ }
115+
116+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
117+ a := &db.Annotation{
118+ URI: event.URI,
119+ AuthorDID: event.DID,
120+ FeedURL: rec.FeedURL,
121+ ArticleURL: rec.ArticleURL,
122+ Quote: db.NullStr(rec.Quote),
123+ Note: db.NullStr(rec.Note),
124+ Tags: db.NullStrTags(rec.Tags),
125+ CreatedAt: sql.NullTime{Time: t, Valid: true},
126+ CID: sql.NullString{String: event.CID, Valid: event.CID != ""},
127+ }
128+ if rec.Rating > 0 {
129+ a.Rating = sql.NullInt64{Int64: int64(rec.Rating), Valid: true}
130+ }
131+ return h.db.CreateAnnotation(ctx, a)
132+
133+ case "delete":
134+ return h.db.DeleteAnnotation(ctx, event.URI)
135+ }
136+ return nil
137+}
modified internal/atproto/lexicon.go +21 -1
@@ -1,6 +1,9 @@
11 package atproto
22
3-import "time"
3+import (
4+ "strings"
5+ "time"
6+)
47
58 type SubscriptionRecord struct {
69 CreatedAt string `json:"createdAt"`
@@ -124,3 +127,20 @@ type GetRecommendationsResponse struct {
124127 Feeds []RecommendedFeed `json:"feeds"`
125128 People []RecommendedPerson `json:"people"`
126129 }
130+
131+type RecordURI struct {
132+ DID string
133+ Collection string
134+ RKey string
135+}
136+
137+func ParseRecordURI(uri string) (RecordURI, bool) {
138+ if !strings.HasPrefix(uri, "at://") {
139+ return RecordURI{}, false
140+ }
141+ parts := strings.SplitN(uri[5:], "/", 3)
142+ if len(parts) != 3 {
143+ return RecordURI{}, false
144+ }
145+ return RecordURI{DID: parts[0], Collection: parts[1], RKey: parts[2]}, true
146+}
@@ -1,6 +1,9 @@
1 package atproto1 package atproto
2 2
3-import "time"3+import (
4+ "strings"
5+ "time"
6+)
4 7
5 type SubscriptionRecord struct {8 type SubscriptionRecord struct {
6 CreatedAt string `json:"createdAt"`9 CreatedAt string `json:"createdAt"`
@@ -124,3 +127,20 @@ type GetRecommendationsResponse struct {
124 Feeds []RecommendedFeed `json:"feeds"`127 Feeds []RecommendedFeed `json:"feeds"`
125 People []RecommendedPerson `json:"people"`128 People []RecommendedPerson `json:"people"`
126 }129 }
130+
131+type RecordURI struct {
132+ DID string
133+ Collection string
134+ RKey string
135+}
136+
137+func ParseRecordURI(uri string) (RecordURI, bool) {
138+ if !strings.HasPrefix(uri, "at://") {
139+ return RecordURI{}, false
140+ }
141+ parts := strings.SplitN(uri[5:], "/", 3)
142+ if len(parts) != 3 {
143+ return RecordURI{}, false
144+ }
145+ return RecordURI{DID: parts[0], Collection: parts[1], RKey: parts[2]}, true
146+}
added internal/atproto/sync.go +150 -0
new file mode 100644
@@ -0,0 +1,150 @@
1+package atproto
2+
3+import (
4+ "context"
5+ "encoding/json"
6+ "log/slog"
7+ "time"
8+
9+ "pkg.rbrt.fr/glean/internal/db"
10+)
11+
12+type Sync struct {
13+ db *db.DB
14+ client *Client
15+ logger *slog.Logger
16+}
17+
18+func NewSync(database *db.DB, client *Client, logger *slog.Logger) *Sync {
19+ return &Sync{db: database, client: client, logger: logger}
20+}
21+
22+func (s *Sync) Run(ctx context.Context, userDID string) error {
23+ s.logger.Info("syncing from PDS", "did", userDID)
24+
25+ if err := s.syncCollection(ctx, userDID, "at.glean.subscription", s.reconcileSubscription); err != nil {
26+ s.logger.Error("sync subscriptions failed", "error", err, "did", userDID)
27+ }
28+ if err := s.syncCollection(ctx, userDID, "at.glean.like", s.reconcileLike); err != nil {
29+ s.logger.Error("sync likes failed", "error", err, "did", userDID)
30+ }
31+ if err := s.syncCollection(ctx, userDID, "at.glean.annotation", s.reconcileAnnotation); err != nil {
32+ s.logger.Error("sync annotations failed", "error", err, "did", userDID)
33+ }
34+
35+ return nil
36+}
37+
38+type reconcileFunc func(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error
39+
40+func (s *Sync) syncCollection(ctx context.Context, userDID, collection string, fn reconcileFunc) error {
41+ cursor := ""
42+ for {
43+ records, next, err := s.client.ListRecords(ctx, userDID, collection, 100, cursor)
44+ if err != nil {
45+ return err
46+ }
47+
48+ for _, r := range records {
49+ if err := fn(ctx, userDID, r.URI, r.CID, r.Value); err != nil {
50+ s.logger.Error("reconcile record error", "error", err, "uri", r.URI)
51+ }
52+ }
53+
54+ if next == "" || len(records) == 0 {
55+ break
56+ }
57+ cursor = next
58+ }
59+ return nil
60+}
61+
62+func (s *Sync) reconcileSubscription(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error {
63+ var rec SubscriptionRecord
64+ if err := json.Unmarshal(value, &rec); err != nil {
65+ return err
66+ }
67+
68+ if rec.FeedURL == "" {
69+ return nil
70+ }
71+
72+ existing, err := s.db.GetSubscription(ctx, userDID, rec.FeedURL)
73+ if err == nil && existing != nil {
74+ if !existing.URI.Valid || existing.URI.String == "" {
75+ return s.db.UpdateSubscriptionURI(ctx, userDID, rec.FeedURL, uri, cid)
76+ }
77+ return nil
78+ }
79+
80+ f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)}
81+ _ = s.db.UpsertFeed(ctx, f)
82+
83+ return s.db.CreateSubscription(ctx, userDID, rec.FeedURL, rec.Category, uri, cid)
84+}
85+
86+func (s *Sync) reconcileLike(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error {
87+ var rec LikeRecord
88+ if err := json.Unmarshal(value, &rec); err != nil {
89+ return err
90+ }
91+
92+ if rec.FeedURL == "" || rec.ArticleURL == "" {
93+ return nil
94+ }
95+
96+ exists, err := s.db.HasLiked(ctx, userDID, rec.FeedURL, rec.ArticleURL)
97+ if err != nil {
98+ return err
99+ }
100+ if exists {
101+ return nil
102+ }
103+
104+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
105+ like := &db.Like{
106+ URI: uri,
107+ AuthorDID: userDID,
108+ FeedURL: rec.FeedURL,
109+ ArticleURL: rec.ArticleURL,
110+ CreatedAt: db.NullTime(t),
111+ CID: db.NullStr(cid),
112+ }
113+ return s.db.CreateLike(ctx, like)
114+}
115+
116+func (s *Sync) reconcileAnnotation(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error {
117+ var rec AnnotationRecord
118+ if err := json.Unmarshal(value, &rec); err != nil {
119+ return err
120+ }
121+
122+ if rec.FeedURL == "" || rec.ArticleURL == "" {
123+ return nil
124+ }
125+
126+ var existing []*db.Annotation
127+ existing, _ = s.db.ListAnnotations(ctx, rec.FeedURL, rec.ArticleURL, userDID, 100, 0)
128+ for _, a := range existing {
129+ if a.URI == uri {
130+ return nil
131+ }
132+ }
133+
134+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
135+ a := &db.Annotation{
136+ URI: uri,
137+ AuthorDID: userDID,
138+ FeedURL: rec.FeedURL,
139+ ArticleURL: rec.ArticleURL,
140+ Quote: db.NullStr(rec.Quote),
141+ Note: db.NullStr(rec.Note),
142+ Tags: db.NullStrTags(rec.Tags),
143+ CreatedAt: db.NullTime(t),
144+ CID: db.NullStr(cid),
145+ }
146+ if rec.Rating > 0 {
147+ a.Rating = db.NullInt(int64(rec.Rating))
148+ }
149+ return s.db.CreateAnnotation(ctx, a)
150+}
new file mode 100644
@@ -0,0 +1,150 @@
1+package atproto
2+
3+import (
4+ "context"
5+ "encoding/json"
6+ "log/slog"
7+ "time"
8+
9+ "pkg.rbrt.fr/glean/internal/db"
10+)
11+
12+type Sync struct {
13+ db *db.DB
14+ client *Client
15+ logger *slog.Logger
16+}
17+
18+func NewSync(database *db.DB, client *Client, logger *slog.Logger) *Sync {
19+ return &Sync{db: database, client: client, logger: logger}
20+}
21+
22+func (s *Sync) Run(ctx context.Context, userDID string) error {
23+ s.logger.Info("syncing from PDS", "did", userDID)
24+
25+ if err := s.syncCollection(ctx, userDID, "at.glean.subscription", s.reconcileSubscription); err != nil {
26+ s.logger.Error("sync subscriptions failed", "error", err, "did", userDID)
27+ }
28+ if err := s.syncCollection(ctx, userDID, "at.glean.like", s.reconcileLike); err != nil {
29+ s.logger.Error("sync likes failed", "error", err, "did", userDID)
30+ }
31+ if err := s.syncCollection(ctx, userDID, "at.glean.annotation", s.reconcileAnnotation); err != nil {
32+ s.logger.Error("sync annotations failed", "error", err, "did", userDID)
33+ }
34+
35+ return nil
36+}
37+
38+type reconcileFunc func(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error
39+
40+func (s *Sync) syncCollection(ctx context.Context, userDID, collection string, fn reconcileFunc) error {
41+ cursor := ""
42+ for {
43+ records, next, err := s.client.ListRecords(ctx, userDID, collection, 100, cursor)
44+ if err != nil {
45+ return err
46+ }
47+
48+ for _, r := range records {
49+ if err := fn(ctx, userDID, r.URI, r.CID, r.Value); err != nil {
50+ s.logger.Error("reconcile record error", "error", err, "uri", r.URI)
51+ }
52+ }
53+
54+ if next == "" || len(records) == 0 {
55+ break
56+ }
57+ cursor = next
58+ }
59+ return nil
60+}
61+
62+func (s *Sync) reconcileSubscription(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error {
63+ var rec SubscriptionRecord
64+ if err := json.Unmarshal(value, &rec); err != nil {
65+ return err
66+ }
67+
68+ if rec.FeedURL == "" {
69+ return nil
70+ }
71+
72+ existing, err := s.db.GetSubscription(ctx, userDID, rec.FeedURL)
73+ if err == nil && existing != nil {
74+ if !existing.URI.Valid || existing.URI.String == "" {
75+ return s.db.UpdateSubscriptionURI(ctx, userDID, rec.FeedURL, uri, cid)
76+ }
77+ return nil
78+ }
79+
80+ f := &db.Feed{FeedURL: rec.FeedURL, Title: db.NullStr(rec.Title)}
81+ _ = s.db.UpsertFeed(ctx, f)
82+
83+ return s.db.CreateSubscription(ctx, userDID, rec.FeedURL, rec.Category, uri, cid)
84+}
85+
86+func (s *Sync) reconcileLike(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error {
87+ var rec LikeRecord
88+ if err := json.Unmarshal(value, &rec); err != nil {
89+ return err
90+ }
91+
92+ if rec.FeedURL == "" || rec.ArticleURL == "" {
93+ return nil
94+ }
95+
96+ exists, err := s.db.HasLiked(ctx, userDID, rec.FeedURL, rec.ArticleURL)
97+ if err != nil {
98+ return err
99+ }
100+ if exists {
101+ return nil
102+ }
103+
104+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
105+ like := &db.Like{
106+ URI: uri,
107+ AuthorDID: userDID,
108+ FeedURL: rec.FeedURL,
109+ ArticleURL: rec.ArticleURL,
110+ CreatedAt: db.NullTime(t),
111+ CID: db.NullStr(cid),
112+ }
113+ return s.db.CreateLike(ctx, like)
114+}
115+
116+func (s *Sync) reconcileAnnotation(ctx context.Context, userDID, uri, cid string, value json.RawMessage) error {
117+ var rec AnnotationRecord
118+ if err := json.Unmarshal(value, &rec); err != nil {
119+ return err
120+ }
121+
122+ if rec.FeedURL == "" || rec.ArticleURL == "" {
123+ return nil
124+ }
125+
126+ var existing []*db.Annotation
127+ existing, _ = s.db.ListAnnotations(ctx, rec.FeedURL, rec.ArticleURL, userDID, 100, 0)
128+ for _, a := range existing {
129+ if a.URI == uri {
130+ return nil
131+ }
132+ }
133+
134+ t, _ := time.Parse(time.RFC3339, rec.CreatedAt)
135+ a := &db.Annotation{
136+ URI: uri,
137+ AuthorDID: userDID,
138+ FeedURL: rec.FeedURL,
139+ ArticleURL: rec.ArticleURL,
140+ Quote: db.NullStr(rec.Quote),
141+ Note: db.NullStr(rec.Note),
142+ Tags: db.NullStrTags(rec.Tags),
143+ CreatedAt: db.NullTime(t),
144+ CID: db.NullStr(cid),
145+ }
146+ if rec.Rating > 0 {
147+ a.Rating = db.NullInt(int64(rec.Rating))
148+ }
149+ return s.db.CreateAnnotation(ctx, a)
150+}
modified internal/db/db.go +160 -134
@@ -2,9 +2,30 @@ package db
22
33 import (
44 "database/sql"
5+ "strings"
6+ "time"
57 _ "github.com/mattn/go-sqlite3"
68 )
79
10+func NullStr(s string) sql.NullString {
11+ return sql.NullString{String: s, Valid: s != ""}
12+}
13+
14+func NullTime(t time.Time) sql.NullTime {
15+ return sql.NullTime{Time: t, Valid: !t.IsZero()}
16+}
17+
18+func NullInt(n int64) sql.NullInt64 {
19+ return sql.NullInt64{Int64: n, Valid: true}
20+}
21+
22+func NullStrTags(tags []string) sql.NullString {
23+ if len(tags) == 0 {
24+ return sql.NullString{}
25+ }
26+ return sql.NullString{String: strings.Join(tags, ","), Valid: true}
27+}
28+
829 type DB struct {
930 *sql.DB
1031 }
@@ -17,7 +38,7 @@ func Open(path string) (*DB, error) {
1738
1839 db.SetMaxOpenConns(1)
1940
20- if err := migrate(db); err != nil {
41+ if err := initSchema(db); err != nil {
2142 db.Close()
2243 return nil, err
2344 }
@@ -25,141 +46,146 @@ func Open(path string) (*DB, error) {
2546 return &DB{db}, nil
2647 }
2748
28-func migrate(db *sql.DB) error {
29- tx, err := db.Begin()
30- if err != nil {
31- return err
32- }
33- defer func() { _ = tx.Rollback() }()
34-
35- stmts := []string{
36- `CREATE TABLE IF NOT EXISTS users (
37- did TEXT PRIMARY KEY,
38- handle TEXT NOT NULL,
39- display_name TEXT,
40- avatar_url TEXT,
41- indexed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
42- updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
43- )`,
44- `CREATE TABLE IF NOT EXISTS feeds (
45- feed_url TEXT PRIMARY KEY,
46- title TEXT,
47- site_url TEXT,
48- description TEXT,
49- feed_type TEXT CHECK(feed_type IN ('rss', 'atom', 'json')),
50- last_fetched_at DATETIME,
51- last_error TEXT,
52- subscriber_count INTEGER NOT NULL DEFAULT 0,
53- etag TEXT,
54- last_modified TEXT,
55- fetch_interval_minutes INTEGER NOT NULL DEFAULT 30,
56- next_fetch_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
57- consecutive_empty_fetches INTEGER NOT NULL DEFAULT 0,
58- error_count INTEGER NOT NULL DEFAULT 0,
59- favicon_url TEXT
60- )`,
61- `CREATE TABLE IF NOT EXISTS subscriptions (
62- id INTEGER PRIMARY KEY AUTOINCREMENT,
63- user_did TEXT NOT NULL REFERENCES users(did),
64- feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
65- category TEXT,
66- added_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
67- UNIQUE(user_did, feed_url)
68- )`,
69- `CREATE TABLE IF NOT EXISTS articles (
70- id INTEGER PRIMARY KEY AUTOINCREMENT,
71- feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
72- guid TEXT NOT NULL,
73- title TEXT NOT NULL DEFAULT '',
74- url TEXT,
75- author TEXT,
76- summary TEXT,
77- content TEXT,
78- published DATETIME,
79- updated DATETIME,
80- fetched_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
81- UNIQUE(feed_url, guid)
82- )`,
83- `CREATE TABLE IF NOT EXISTS read_state (
84- user_did TEXT NOT NULL REFERENCES users(did),
85- article_id INTEGER NOT NULL REFERENCES articles(id),
86- is_read BOOLEAN NOT NULL DEFAULT 0,
87- read_at DATETIME,
88- is_starred BOOLEAN NOT NULL DEFAULT 0,
89- starred_at DATETIME,
90- PRIMARY KEY (user_did, article_id)
91- )`,
92- `CREATE TABLE IF NOT EXISTS annotations (
93- id INTEGER PRIMARY KEY AUTOINCREMENT,
94- uri TEXT NOT NULL UNIQUE,
95- author_did TEXT NOT NULL REFERENCES users(did),
96- feed_url TEXT NOT NULL,
97- article_url TEXT NOT NULL,
98- quote TEXT,
99- note TEXT,
100- tags TEXT,
101- rating INTEGER,
102- created_at DATETIME NOT NULL,
103- cid TEXT
104- )`,
105- `CREATE TABLE IF NOT EXISTS likes (
106- id INTEGER PRIMARY KEY AUTOINCREMENT,
107- uri TEXT NOT NULL UNIQUE,
108- author_did TEXT NOT NULL REFERENCES users(did),
109- feed_url TEXT NOT NULL,
110- article_url TEXT NOT NULL,
111- created_at DATETIME NOT NULL,
112- cid TEXT,
113- UNIQUE(author_did, feed_url, article_url)
114- )`,
115- `CREATE TABLE IF NOT EXISTS feed_similarity (
116- feed_a TEXT NOT NULL REFERENCES feeds(feed_url),
117- feed_b TEXT NOT NULL REFERENCES feeds(feed_url),
118- jaccard REAL NOT NULL,
119- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
120- PRIMARY KEY (feed_a, feed_b),
121- CHECK(feed_a < feed_b)
122- )`,
123- `CREATE TABLE IF NOT EXISTS user_similarity (
124- user_a TEXT NOT NULL REFERENCES users(did),
125- user_b TEXT NOT NULL REFERENCES users(did),
126- jaccard REAL NOT NULL,
127- common_feeds INTEGER NOT NULL,
128- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
129- PRIMARY KEY (user_a, user_b),
130- CHECK(user_a < user_b)
131- )`,
132- `CREATE TABLE IF NOT EXISTS user_feed_recommendations (
133- user_did TEXT NOT NULL REFERENCES users(did),
134- feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
135- score REAL NOT NULL,
136- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
137- PRIMARY KEY (user_did, feed_url)
138- )`,
139- `CREATE TABLE IF NOT EXISTS user_article_recommendations (
140- user_did TEXT NOT NULL REFERENCES users(did),
141- feed_url TEXT NOT NULL,
142- article_url TEXT NOT NULL,
143- score REAL NOT NULL,
144- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
145- PRIMARY KEY (user_did, feed_url, article_url)
146- )`,
147- `CREATE INDEX IF NOT EXISTS idx_subscriptions_feed ON subscriptions(feed_url)`,
148- `CREATE INDEX IF NOT EXISTS idx_subscriptions_user ON subscriptions(user_did)`,
149- `CREATE INDEX IF NOT EXISTS idx_articles_feed ON articles(feed_url)`,
150- `CREATE INDEX IF NOT EXISTS idx_articles_published ON articles(published DESC)`,
151- `CREATE INDEX IF NOT EXISTS idx_read_state_unread ON read_state(user_did, is_read) WHERE is_read = 0`,
152- `CREATE INDEX IF NOT EXISTS idx_read_state_starred ON read_state(user_did, is_starred) WHERE is_starred = 1`,
153- `CREATE INDEX IF NOT EXISTS idx_annotations_article ON annotations(article_url)`,
154- `CREATE INDEX IF NOT EXISTS idx_likes_article ON likes(feed_url, article_url)`,
155- `CREATE INDEX IF NOT EXISTS idx_likes_author ON likes(author_did)`,
156- }
157-
158- for _, s := range stmts {
159- if _, err := tx.Exec(s); err != nil {
49+func initSchema(db *sql.DB) error {
50+ for _, s := range schema {
51+ if _, err := db.Exec(s); err != nil {
16052 return err
16153 }
16254 }
55+ return nil
56+}
16357
164- return tx.Commit()
58+var schema = []string{
59+ `CREATE TABLE IF NOT EXISTS users (
60+ did TEXT PRIMARY KEY,
61+ handle TEXT NOT NULL,
62+ display_name TEXT,
63+ avatar_url TEXT,
64+ indexed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
65+ updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
66+ )`,
67+ `CREATE TABLE IF NOT EXISTS feeds (
68+ feed_url TEXT PRIMARY KEY,
69+ title TEXT,
70+ site_url TEXT,
71+ description TEXT,
72+ feed_type TEXT CHECK(feed_type IN ('rss', 'atom', 'json')),
73+ last_fetched_at DATETIME,
74+ last_error TEXT,
75+ subscriber_count INTEGER NOT NULL DEFAULT 0,
76+ etag TEXT,
77+ last_modified TEXT,
78+ fetch_interval_minutes INTEGER NOT NULL DEFAULT 30,
79+ next_fetch_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
80+ consecutive_empty_fetches INTEGER NOT NULL DEFAULT 0,
81+ error_count INTEGER NOT NULL DEFAULT 0,
82+ favicon_url TEXT
83+ )`,
84+ `CREATE TABLE IF NOT EXISTS subscriptions (
85+ id INTEGER PRIMARY KEY AUTOINCREMENT,
86+ user_did TEXT NOT NULL REFERENCES users(did),
87+ feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
88+ category TEXT,
89+ added_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
90+ uri TEXT,
91+ cid TEXT,
92+ UNIQUE(user_did, feed_url)
93+ )`,
94+ `CREATE TABLE IF NOT EXISTS articles (
95+ id INTEGER PRIMARY KEY AUTOINCREMENT,
96+ feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
97+ guid TEXT NOT NULL,
98+ title TEXT NOT NULL DEFAULT '',
99+ url TEXT,
100+ author TEXT,
101+ summary TEXT,
102+ content TEXT,
103+ published DATETIME,
104+ updated DATETIME,
105+ fetched_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
106+ UNIQUE(feed_url, guid)
107+ )`,
108+ `CREATE TABLE IF NOT EXISTS read_state (
109+ user_did TEXT NOT NULL REFERENCES users(did),
110+ article_id INTEGER NOT NULL REFERENCES articles(id),
111+ is_read BOOLEAN NOT NULL DEFAULT 0,
112+ read_at DATETIME,
113+ is_starred BOOLEAN NOT NULL DEFAULT 0,
114+ starred_at DATETIME,
115+ PRIMARY KEY (user_did, article_id)
116+ )`,
117+ `CREATE TABLE IF NOT EXISTS annotations (
118+ id INTEGER PRIMARY KEY AUTOINCREMENT,
119+ uri TEXT NOT NULL UNIQUE,
120+ author_did TEXT NOT NULL REFERENCES users(did),
121+ feed_url TEXT NOT NULL,
122+ article_url TEXT NOT NULL,
123+ quote TEXT,
124+ note TEXT,
125+ tags TEXT,
126+ rating INTEGER,
127+ created_at DATETIME NOT NULL,
128+ cid TEXT
129+ )`,
130+ `CREATE TABLE IF NOT EXISTS likes (
131+ id INTEGER PRIMARY KEY AUTOINCREMENT,
132+ uri TEXT NOT NULL UNIQUE,
133+ author_did TEXT NOT NULL REFERENCES users(did),
134+ feed_url TEXT NOT NULL,
135+ article_url TEXT NOT NULL,
136+ created_at DATETIME NOT NULL,
137+ cid TEXT,
138+ UNIQUE(author_did, feed_url, article_url)
139+ )`,
140+ `CREATE TABLE IF NOT EXISTS feed_similarity (
141+ feed_a TEXT NOT NULL REFERENCES feeds(feed_url),
142+ feed_b TEXT NOT NULL REFERENCES feeds(feed_url),
143+ jaccard REAL NOT NULL,
144+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
145+ PRIMARY KEY (feed_a, feed_b),
146+ CHECK(feed_a < feed_b)
147+ )`,
148+ `CREATE TABLE IF NOT EXISTS user_similarity (
149+ user_a TEXT NOT NULL REFERENCES users(did),
150+ user_b TEXT NOT NULL REFERENCES users(did),
151+ jaccard REAL NOT NULL,
152+ common_feeds INTEGER NOT NULL,
153+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
154+ PRIMARY KEY (user_a, user_b),
155+ CHECK(user_a < user_b)
156+ )`,
157+ `CREATE TABLE IF NOT EXISTS user_feed_recommendations (
158+ user_did TEXT NOT NULL REFERENCES users(did),
159+ feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
160+ score REAL NOT NULL,
161+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
162+ PRIMARY KEY (user_did, feed_url)
163+ )`,
164+ `CREATE TABLE IF NOT EXISTS user_article_recommendations (
165+ user_did TEXT NOT NULL REFERENCES users(did),
166+ feed_url TEXT NOT NULL,
167+ article_url TEXT NOT NULL,
168+ score REAL NOT NULL,
169+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
170+ PRIMARY KEY (user_did, feed_url, article_url)
171+ )`,
172+ `CREATE TABLE IF NOT EXISTS oauth_auth_requests (
173+ state TEXT PRIMARY KEY,
174+ data TEXT NOT NULL
175+ )`,
176+ `CREATE TABLE IF NOT EXISTS oauth_sessions (
177+ account_did TEXT NOT NULL,
178+ session_id TEXT NOT NULL,
179+ data TEXT NOT NULL,
180+ PRIMARY KEY (account_did, session_id)
181+ )`,
182+ `CREATE INDEX IF NOT EXISTS idx_subscriptions_feed ON subscriptions(feed_url)`,
183+ `CREATE INDEX IF NOT EXISTS idx_subscriptions_user ON subscriptions(user_did)`,
184+ `CREATE INDEX IF NOT EXISTS idx_articles_feed ON articles(feed_url)`,
185+ `CREATE INDEX IF NOT EXISTS idx_articles_published ON articles(published DESC)`,
186+ `CREATE INDEX IF NOT EXISTS idx_read_state_unread ON read_state(user_did, is_read) WHERE is_read = 0`,
187+ `CREATE INDEX IF NOT EXISTS idx_read_state_starred ON read_state(user_did, is_starred) WHERE is_starred = 1`,
188+ `CREATE INDEX IF NOT EXISTS idx_annotations_article ON annotations(article_url)`,
189+ `CREATE INDEX IF NOT EXISTS idx_likes_article ON likes(feed_url, article_url)`,
190+ `CREATE INDEX IF NOT EXISTS idx_likes_author ON likes(author_did)`,
165191 }
@@ -2,9 +2,30 @@ package db
2 2
3 import (3 import (
4 "database/sql"4 "database/sql"
5+ "strings"
6+ "time"
5 _ "github.com/mattn/go-sqlite3"7 _ "github.com/mattn/go-sqlite3"
6 )8 )
7 9
10+func NullStr(s string) sql.NullString {
11+ return sql.NullString{String: s, Valid: s != ""}
12+}
13+
14+func NullTime(t time.Time) sql.NullTime {
15+ return sql.NullTime{Time: t, Valid: !t.IsZero()}
16+}
17+
18+func NullInt(n int64) sql.NullInt64 {
19+ return sql.NullInt64{Int64: n, Valid: true}
20+}
21+
22+func NullStrTags(tags []string) sql.NullString {
23+ if len(tags) == 0 {
24+ return sql.NullString{}
25+ }
26+ return sql.NullString{String: strings.Join(tags, ","), Valid: true}
27+}
28+
8 type DB struct {29 type DB struct {
9 *sql.DB30 *sql.DB
10 }31 }
@@ -17,7 +38,7 @@ func Open(path string) (*DB, error) {
17 38
18 db.SetMaxOpenConns(1)39 db.SetMaxOpenConns(1)
19 40
20- if err := migrate(db); err != nil {41+ if err := initSchema(db); err != nil {
21 db.Close()42 db.Close()
22 return nil, err43 return nil, err
23 }44 }
@@ -25,141 +46,146 @@ func Open(path string) (*DB, error) {
25 return &DB{db}, nil46 return &DB{db}, nil
26 }47 }
27 48
28-func migrate(db *sql.DB) error {49+func initSchema(db *sql.DB) error {
29- tx, err := db.Begin()50+ for _, s := range schema {
30- if err != nil {51+ if _, err := db.Exec(s); err != nil {
31- return err
32- }
33- defer func() { _ = tx.Rollback() }()
34-
35- stmts := []string{
36- `CREATE TABLE IF NOT EXISTS users (
37- did TEXT PRIMARY KEY,
38- handle TEXT NOT NULL,
39- display_name TEXT,
40- avatar_url TEXT,
41- indexed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
42- updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
43- )`,
44- `CREATE TABLE IF NOT EXISTS feeds (
45- feed_url TEXT PRIMARY KEY,
46- title TEXT,
47- site_url TEXT,
48- description TEXT,
49- feed_type TEXT CHECK(feed_type IN ('rss', 'atom', 'json')),
50- last_fetched_at DATETIME,
51- last_error TEXT,
52- subscriber_count INTEGER NOT NULL DEFAULT 0,
53- etag TEXT,
54- last_modified TEXT,
55- fetch_interval_minutes INTEGER NOT NULL DEFAULT 30,
56- next_fetch_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
57- consecutive_empty_fetches INTEGER NOT NULL DEFAULT 0,
58- error_count INTEGER NOT NULL DEFAULT 0,
59- favicon_url TEXT
60- )`,
61- `CREATE TABLE IF NOT EXISTS subscriptions (
62- id INTEGER PRIMARY KEY AUTOINCREMENT,
63- user_did TEXT NOT NULL REFERENCES users(did),
64- feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
65- category TEXT,
66- added_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
67- UNIQUE(user_did, feed_url)
68- )`,
69- `CREATE TABLE IF NOT EXISTS articles (
70- id INTEGER PRIMARY KEY AUTOINCREMENT,
71- feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
72- guid TEXT NOT NULL,
73- title TEXT NOT NULL DEFAULT '',
74- url TEXT,
75- author TEXT,
76- summary TEXT,
77- content TEXT,
78- published DATETIME,
79- updated DATETIME,
80- fetched_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
81- UNIQUE(feed_url, guid)
82- )`,
83- `CREATE TABLE IF NOT EXISTS read_state (
84- user_did TEXT NOT NULL REFERENCES users(did),
85- article_id INTEGER NOT NULL REFERENCES articles(id),
86- is_read BOOLEAN NOT NULL DEFAULT 0,
87- read_at DATETIME,
88- is_starred BOOLEAN NOT NULL DEFAULT 0,
89- starred_at DATETIME,
90- PRIMARY KEY (user_did, article_id)
91- )`,
92- `CREATE TABLE IF NOT EXISTS annotations (
93- id INTEGER PRIMARY KEY AUTOINCREMENT,
94- uri TEXT NOT NULL UNIQUE,
95- author_did TEXT NOT NULL REFERENCES users(did),
96- feed_url TEXT NOT NULL,
97- article_url TEXT NOT NULL,
98- quote TEXT,
99- note TEXT,
100- tags TEXT,
101- rating INTEGER,
102- created_at DATETIME NOT NULL,
103- cid TEXT
104- )`,
105- `CREATE TABLE IF NOT EXISTS likes (
106- id INTEGER PRIMARY KEY AUTOINCREMENT,
107- uri TEXT NOT NULL UNIQUE,
108- author_did TEXT NOT NULL REFERENCES users(did),
109- feed_url TEXT NOT NULL,
110- article_url TEXT NOT NULL,
111- created_at DATETIME NOT NULL,
112- cid TEXT,
113- UNIQUE(author_did, feed_url, article_url)
114- )`,
115- `CREATE TABLE IF NOT EXISTS feed_similarity (
116- feed_a TEXT NOT NULL REFERENCES feeds(feed_url),
117- feed_b TEXT NOT NULL REFERENCES feeds(feed_url),
118- jaccard REAL NOT NULL,
119- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
120- PRIMARY KEY (feed_a, feed_b),
121- CHECK(feed_a < feed_b)
122- )`,
123- `CREATE TABLE IF NOT EXISTS user_similarity (
124- user_a TEXT NOT NULL REFERENCES users(did),
125- user_b TEXT NOT NULL REFERENCES users(did),
126- jaccard REAL NOT NULL,
127- common_feeds INTEGER NOT NULL,
128- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
129- PRIMARY KEY (user_a, user_b),
130- CHECK(user_a < user_b)
131- )`,
132- `CREATE TABLE IF NOT EXISTS user_feed_recommendations (
133- user_did TEXT NOT NULL REFERENCES users(did),
134- feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
135- score REAL NOT NULL,
136- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
137- PRIMARY KEY (user_did, feed_url)
138- )`,
139- `CREATE TABLE IF NOT EXISTS user_article_recommendations (
140- user_did TEXT NOT NULL REFERENCES users(did),
141- feed_url TEXT NOT NULL,
142- article_url TEXT NOT NULL,
143- score REAL NOT NULL,
144- computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
145- PRIMARY KEY (user_did, feed_url, article_url)
146- )`,
147- `CREATE INDEX IF NOT EXISTS idx_subscriptions_feed ON subscriptions(feed_url)`,
148- `CREATE INDEX IF NOT EXISTS idx_subscriptions_user ON subscriptions(user_did)`,
149- `CREATE INDEX IF NOT EXISTS idx_articles_feed ON articles(feed_url)`,
150- `CREATE INDEX IF NOT EXISTS idx_articles_published ON articles(published DESC)`,
151- `CREATE INDEX IF NOT EXISTS idx_read_state_unread ON read_state(user_did, is_read) WHERE is_read = 0`,
152- `CREATE INDEX IF NOT EXISTS idx_read_state_starred ON read_state(user_did, is_starred) WHERE is_starred = 1`,
153- `CREATE INDEX IF NOT EXISTS idx_annotations_article ON annotations(article_url)`,
154- `CREATE INDEX IF NOT EXISTS idx_likes_article ON likes(feed_url, article_url)`,
155- `CREATE INDEX IF NOT EXISTS idx_likes_author ON likes(author_did)`,
156- }
157-
158- for _, s := range stmts {
159- if _, err := tx.Exec(s); err != nil {
160 return err52 return err
161 }53 }
162 }54 }
55+ return nil
56+}
163 57
164- return tx.Commit()58+var schema = []string{
59+ `CREATE TABLE IF NOT EXISTS users (
60+ did TEXT PRIMARY KEY,
61+ handle TEXT NOT NULL,
62+ display_name TEXT,
63+ avatar_url TEXT,
64+ indexed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
65+ updated_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP
66+ )`,
67+ `CREATE TABLE IF NOT EXISTS feeds (
68+ feed_url TEXT PRIMARY KEY,
69+ title TEXT,
70+ site_url TEXT,
71+ description TEXT,
72+ feed_type TEXT CHECK(feed_type IN ('rss', 'atom', 'json')),
73+ last_fetched_at DATETIME,
74+ last_error TEXT,
75+ subscriber_count INTEGER NOT NULL DEFAULT 0,
76+ etag TEXT,
77+ last_modified TEXT,
78+ fetch_interval_minutes INTEGER NOT NULL DEFAULT 30,
79+ next_fetch_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
80+ consecutive_empty_fetches INTEGER NOT NULL DEFAULT 0,
81+ error_count INTEGER NOT NULL DEFAULT 0,
82+ favicon_url TEXT
83+ )`,
84+ `CREATE TABLE IF NOT EXISTS subscriptions (
85+ id INTEGER PRIMARY KEY AUTOINCREMENT,
86+ user_did TEXT NOT NULL REFERENCES users(did),
87+ feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
88+ category TEXT,
89+ added_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
90+ uri TEXT,
91+ cid TEXT,
92+ UNIQUE(user_did, feed_url)
93+ )`,
94+ `CREATE TABLE IF NOT EXISTS articles (
95+ id INTEGER PRIMARY KEY AUTOINCREMENT,
96+ feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
97+ guid TEXT NOT NULL,
98+ title TEXT NOT NULL DEFAULT '',
99+ url TEXT,
100+ author TEXT,
101+ summary TEXT,
102+ content TEXT,
103+ published DATETIME,
104+ updated DATETIME,
105+ fetched_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
106+ UNIQUE(feed_url, guid)
107+ )`,
108+ `CREATE TABLE IF NOT EXISTS read_state (
109+ user_did TEXT NOT NULL REFERENCES users(did),
110+ article_id INTEGER NOT NULL REFERENCES articles(id),
111+ is_read BOOLEAN NOT NULL DEFAULT 0,
112+ read_at DATETIME,
113+ is_starred BOOLEAN NOT NULL DEFAULT 0,
114+ starred_at DATETIME,
115+ PRIMARY KEY (user_did, article_id)
116+ )`,
117+ `CREATE TABLE IF NOT EXISTS annotations (
118+ id INTEGER PRIMARY KEY AUTOINCREMENT,
119+ uri TEXT NOT NULL UNIQUE,
120+ author_did TEXT NOT NULL REFERENCES users(did),
121+ feed_url TEXT NOT NULL,
122+ article_url TEXT NOT NULL,
123+ quote TEXT,
124+ note TEXT,
125+ tags TEXT,
126+ rating INTEGER,
127+ created_at DATETIME NOT NULL,
128+ cid TEXT
129+ )`,
130+ `CREATE TABLE IF NOT EXISTS likes (
131+ id INTEGER PRIMARY KEY AUTOINCREMENT,
132+ uri TEXT NOT NULL UNIQUE,
133+ author_did TEXT NOT NULL REFERENCES users(did),
134+ feed_url TEXT NOT NULL,
135+ article_url TEXT NOT NULL,
136+ created_at DATETIME NOT NULL,
137+ cid TEXT,
138+ UNIQUE(author_did, feed_url, article_url)
139+ )`,
140+ `CREATE TABLE IF NOT EXISTS feed_similarity (
141+ feed_a TEXT NOT NULL REFERENCES feeds(feed_url),
142+ feed_b TEXT NOT NULL REFERENCES feeds(feed_url),
143+ jaccard REAL NOT NULL,
144+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
145+ PRIMARY KEY (feed_a, feed_b),
146+ CHECK(feed_a < feed_b)
147+ )`,
148+ `CREATE TABLE IF NOT EXISTS user_similarity (
149+ user_a TEXT NOT NULL REFERENCES users(did),
150+ user_b TEXT NOT NULL REFERENCES users(did),
151+ jaccard REAL NOT NULL,
152+ common_feeds INTEGER NOT NULL,
153+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
154+ PRIMARY KEY (user_a, user_b),
155+ CHECK(user_a < user_b)
156+ )`,
157+ `CREATE TABLE IF NOT EXISTS user_feed_recommendations (
158+ user_did TEXT NOT NULL REFERENCES users(did),
159+ feed_url TEXT NOT NULL REFERENCES feeds(feed_url),
160+ score REAL NOT NULL,
161+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
162+ PRIMARY KEY (user_did, feed_url)
163+ )`,
164+ `CREATE TABLE IF NOT EXISTS user_article_recommendations (
165+ user_did TEXT NOT NULL REFERENCES users(did),
166+ feed_url TEXT NOT NULL,
167+ article_url TEXT NOT NULL,
168+ score REAL NOT NULL,
169+ computed_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP,
170+ PRIMARY KEY (user_did, feed_url, article_url)
171+ )`,
172+ `CREATE TABLE IF NOT EXISTS oauth_auth_requests (
173+ state TEXT PRIMARY KEY,
174+ data TEXT NOT NULL
175+ )`,
176+ `CREATE TABLE IF NOT EXISTS oauth_sessions (
177+ account_did TEXT NOT NULL,
178+ session_id TEXT NOT NULL,
179+ data TEXT NOT NULL,
180+ PRIMARY KEY (account_did, session_id)
181+ )`,
182+ `CREATE INDEX IF NOT EXISTS idx_subscriptions_feed ON subscriptions(feed_url)`,
183+ `CREATE INDEX IF NOT EXISTS idx_subscriptions_user ON subscriptions(user_did)`,
184+ `CREATE INDEX IF NOT EXISTS idx_articles_feed ON articles(feed_url)`,
185+ `CREATE INDEX IF NOT EXISTS idx_articles_published ON articles(published DESC)`,
186+ `CREATE INDEX IF NOT EXISTS idx_read_state_unread ON read_state(user_did, is_read) WHERE is_read = 0`,
187+ `CREATE INDEX IF NOT EXISTS idx_read_state_starred ON read_state(user_did, is_starred) WHERE is_starred = 1`,
188+ `CREATE INDEX IF NOT EXISTS idx_annotations_article ON annotations(article_url)`,
189+ `CREATE INDEX IF NOT EXISTS idx_likes_article ON likes(feed_url, article_url)`,
190+ `CREATE INDEX IF NOT EXISTS idx_likes_author ON likes(author_did)`,
165 }191 }
modified internal/db/feed.go +37 -6
@@ -34,6 +34,8 @@ type Subscription struct {
3434 AddedAt sql.NullTime
3535 UnreadCount int
3636 FetchInterval int
37+ URI sql.NullString
38+ CID sql.NullString
3739 }
3840
3941 func (db *DB) UpsertFeed(ctx context.Context, feed *Feed) error {
@@ -124,17 +126,31 @@ func (db *DB) DecrementSubscriberCount(ctx context.Context, feedURL string) erro
124126 return err
125127 }
126128
127-func (db *DB) CreateSubscription(ctx context.Context, userDID, feedURL, category string) error {
129+func (db *DB) CreateSubscription(ctx context.Context, userDID, feedURL, category, uri, cid string) error {
128130 _, err := db.ExecContext(ctx, `
129- INSERT INTO subscriptions (user_did, feed_url, category)
130- VALUES (?, ?, ?)
131- `, userDID, feedURL, category)
131+ INSERT INTO subscriptions (user_did, feed_url, category, uri, cid)
132+ VALUES (?, ?, ?, ?, ?)
133+ `, userDID, feedURL, category, uriOrNil(category, uri), uriOrNil(category, cid))
132134 if err != nil {
133135 return err
134136 }
135137 return db.IncrementSubscriberCount(ctx, feedURL)
136138 }
137139
140+func (db *DB) UpdateSubscriptionURI(ctx context.Context, userDID, feedURL, uri, cid string) error {
141+ _, err := db.ExecContext(ctx, `
142+ UPDATE subscriptions SET uri = ?, cid = ? WHERE user_did = ? AND feed_url = ?
143+ `, uri, cid, userDID, feedURL)
144+ return err
145+}
146+
147+func uriOrNil(category, v string) any {
148+ if v == "" {
149+ return nil
150+ }
151+ return v
152+}
153+
138154 func (db *DB) DeleteSubscription(ctx context.Context, userDID, feedURL string) error {
139155 _, err := db.ExecContext(ctx, `
140156 DELETE FROM subscriptions WHERE user_did = ? AND feed_url = ?
@@ -145,9 +161,24 @@ func (db *DB) DeleteSubscription(ctx context.Context, userDID, feedURL string) e
145161 return db.DecrementSubscriberCount(ctx, feedURL)
146162 }
147163
164+func (db *DB) GetSubscription(ctx context.Context, userDID, feedURL string) (*Subscription, error) {
165+ s := &Subscription{}
166+ err := db.QueryRowContext(ctx, `
167+ SELECT s.id, s.user_did, s.feed_url, COALESCE(f.title, ''), s.category, s.added_at,
168+ COALESCE(f.fetch_interval_minutes, 30), s.uri, s.cid
169+ FROM subscriptions s
170+ LEFT JOIN feeds f ON s.feed_url = f.feed_url
171+ WHERE s.user_did = ? AND s.feed_url = ?
172+ `, userDID, feedURL).Scan(&s.ID, &s.UserDID, &s.FeedURL, &s.FeedTitle, &s.Category, &s.AddedAt, &s.FetchInterval, &s.URI, &s.CID)
173+ if err != nil {
174+ return nil, err
175+ }
176+ return s, nil
177+}
178+
148179 func (db *DB) ListSubscriptions(ctx context.Context, userDID, category string, limit, offset int) ([]*Subscription, error) {
149180 query := `SELECT s.id, s.user_did, s.feed_url, COALESCE(f.title, ''), s.category, s.added_at,
150- COALESCE(f.fetch_interval_minutes, 30)
181+ COALESCE(f.fetch_interval_minutes, 30), s.uri, s.cid
151182 FROM subscriptions s
152183 LEFT JOIN feeds f ON s.feed_url = f.feed_url
153184 WHERE s.user_did = ?`
@@ -169,7 +200,7 @@ func (db *DB) ListSubscriptions(ctx context.Context, userDID, category string, l
169200 var subs []*Subscription
170201 for rows.Next() {
171202 s := &Subscription{}
172- if err := rows.Scan(&s.ID, &s.UserDID, &s.FeedURL, &s.FeedTitle, &s.Category, &s.AddedAt, &s.FetchInterval); err != nil {
203+ if err := rows.Scan(&s.ID, &s.UserDID, &s.FeedURL, &s.FeedTitle, &s.Category, &s.AddedAt, &s.FetchInterval, &s.URI, &s.CID); err != nil {
173204 return nil, err
174205 }
175206 subs = append(subs, s)
@@ -34,6 +34,8 @@ type Subscription struct {
34 AddedAt sql.NullTime34 AddedAt sql.NullTime
35 UnreadCount int35 UnreadCount int
36 FetchInterval int36 FetchInterval int
37+ URI sql.NullString
38+ CID sql.NullString
37 }39 }
38 40
39 func (db *DB) UpsertFeed(ctx context.Context, feed *Feed) error {41 func (db *DB) UpsertFeed(ctx context.Context, feed *Feed) error {
@@ -124,17 +126,31 @@ func (db *DB) DecrementSubscriberCount(ctx context.Context, feedURL string) erro
124 return err126 return err
125 }127 }
126 128
127-func (db *DB) CreateSubscription(ctx context.Context, userDID, feedURL, category string) error {129+func (db *DB) CreateSubscription(ctx context.Context, userDID, feedURL, category, uri, cid string) error {
128 _, err := db.ExecContext(ctx, `130 _, err := db.ExecContext(ctx, `
129- INSERT INTO subscriptions (user_did, feed_url, category)131+ INSERT INTO subscriptions (user_did, feed_url, category, uri, cid)
130- VALUES (?, ?, ?)132+ VALUES (?, ?, ?, ?, ?)
131- `, userDID, feedURL, category)133+ `, userDID, feedURL, category, uriOrNil(category, uri), uriOrNil(category, cid))
132 if err != nil {134 if err != nil {
133 return err135 return err
134 }136 }
135 return db.IncrementSubscriberCount(ctx, feedURL)137 return db.IncrementSubscriberCount(ctx, feedURL)
136 }138 }
137 139
140+func (db *DB) UpdateSubscriptionURI(ctx context.Context, userDID, feedURL, uri, cid string) error {
141+ _, err := db.ExecContext(ctx, `
142+ UPDATE subscriptions SET uri = ?, cid = ? WHERE user_did = ? AND feed_url = ?
143+ `, uri, cid, userDID, feedURL)
144+ return err
145+}
146+
147+func uriOrNil(category, v string) any {
148+ if v == "" {
149+ return nil
150+ }
151+ return v
152+}
153+
138 func (db *DB) DeleteSubscription(ctx context.Context, userDID, feedURL string) error {154 func (db *DB) DeleteSubscription(ctx context.Context, userDID, feedURL string) error {
139 _, err := db.ExecContext(ctx, `155 _, err := db.ExecContext(ctx, `
140 DELETE FROM subscriptions WHERE user_did = ? AND feed_url = ?156 DELETE FROM subscriptions WHERE user_did = ? AND feed_url = ?
@@ -145,9 +161,24 @@ func (db *DB) DeleteSubscription(ctx context.Context, userDID, feedURL string) e
145 return db.DecrementSubscriberCount(ctx, feedURL)161 return db.DecrementSubscriberCount(ctx, feedURL)
146 }162 }
147 163
164+func (db *DB) GetSubscription(ctx context.Context, userDID, feedURL string) (*Subscription, error) {
165+ s := &Subscription{}
166+ err := db.QueryRowContext(ctx, `
167+ SELECT s.id, s.user_did, s.feed_url, COALESCE(f.title, ''), s.category, s.added_at,
168+ COALESCE(f.fetch_interval_minutes, 30), s.uri, s.cid
169+ FROM subscriptions s
170+ LEFT JOIN feeds f ON s.feed_url = f.feed_url
171+ WHERE s.user_did = ? AND s.feed_url = ?
172+ `, userDID, feedURL).Scan(&s.ID, &s.UserDID, &s.FeedURL, &s.FeedTitle, &s.Category, &s.AddedAt, &s.FetchInterval, &s.URI, &s.CID)
173+ if err != nil {
174+ return nil, err
175+ }
176+ return s, nil
177+}
178+
148 func (db *DB) ListSubscriptions(ctx context.Context, userDID, category string, limit, offset int) ([]*Subscription, error) {179 func (db *DB) ListSubscriptions(ctx context.Context, userDID, category string, limit, offset int) ([]*Subscription, error) {
149 query := `SELECT s.id, s.user_did, s.feed_url, COALESCE(f.title, ''), s.category, s.added_at,180 query := `SELECT s.id, s.user_did, s.feed_url, COALESCE(f.title, ''), s.category, s.added_at,
150- COALESCE(f.fetch_interval_minutes, 30)181+ COALESCE(f.fetch_interval_minutes, 30), s.uri, s.cid
151 FROM subscriptions s182 FROM subscriptions s
152 LEFT JOIN feeds f ON s.feed_url = f.feed_url183 LEFT JOIN feeds f ON s.feed_url = f.feed_url
153 WHERE s.user_did = ?`184 WHERE s.user_did = ?`
@@ -169,7 +200,7 @@ func (db *DB) ListSubscriptions(ctx context.Context, userDID, category string, l
169 var subs []*Subscription200 var subs []*Subscription
170 for rows.Next() {201 for rows.Next() {
171 s := &Subscription{}202 s := &Subscription{}
172- if err := rows.Scan(&s.ID, &s.UserDID, &s.FeedURL, &s.FeedTitle, &s.Category, &s.AddedAt, &s.FetchInterval); err != nil {203+ if err := rows.Scan(&s.ID, &s.UserDID, &s.FeedURL, &s.FeedTitle, &s.Category, &s.AddedAt, &s.FetchInterval, &s.URI, &s.CID); err != nil {
173 return nil, err204 return nil, err
174 }205 }
175 subs = append(subs, s)206 subs = append(subs, s)
modified internal/db/oauth_store.go +20 -21
@@ -18,27 +18,6 @@ func NewOAuthStore(db *DB) *OAuthStore {
1818 return &OAuthStore{db: db}
1919 }
2020
21-func (s *OAuthStore) Init(ctx context.Context) error {
22- stmts := []string{
23- `CREATE TABLE IF NOT EXISTS oauth_auth_requests (
24- state TEXT PRIMARY KEY,
25- data TEXT NOT NULL
26- )`,
27- `CREATE TABLE IF NOT EXISTS oauth_sessions (
28- account_did TEXT NOT NULL,
29- session_id TEXT NOT NULL,
30- data TEXT NOT NULL,
31- PRIMARY KEY (account_did, session_id)
32- )`,
33- }
34- for _, stmt := range stmts {
35- if _, err := s.db.ExecContext(ctx, stmt); err != nil {
36- return err
37- }
38- }
39- return nil
40-}
41-
4221 func (s *OAuthStore) GetSession(ctx context.Context, did syntax.DID, sessionID string) (*oauth.ClientSessionData, error) {
4322 var data []byte
4423 err := s.db.QueryRowContext(ctx, `
@@ -77,6 +56,26 @@ func (s *OAuthStore) DeleteSession(ctx context.Context, did syntax.DID, sessionI
7756 return err
7857 }
7958
59+func (s *OAuthStore) ListSessionsForDID(ctx context.Context, did string) ([]string, error) {
60+ rows, err := s.db.QueryContext(ctx, `
61+ SELECT session_id FROM oauth_sessions WHERE account_did = ? ORDER BY ROWID DESC
62+ `, did)
63+ if err != nil {
64+ return nil, err
65+ }
66+ defer rows.Close()
67+
68+ var ids []string
69+ for rows.Next() {
70+ var id string
71+ if err := rows.Scan(&id); err != nil {
72+ return nil, err
73+ }
74+ ids = append(ids, id)
75+ }
76+ return ids, rows.Err()
77+}
78+
8079 func (s *OAuthStore) GetAuthRequestInfo(ctx context.Context, state string) (*oauth.AuthRequestData, error) {
8180 var data []byte
8281 err := s.db.QueryRowContext(ctx, `
@@ -18,27 +18,6 @@ func NewOAuthStore(db *DB) *OAuthStore {
18 return &OAuthStore{db: db}18 return &OAuthStore{db: db}
19 }19 }
20 20
21-func (s *OAuthStore) Init(ctx context.Context) error {
22- stmts := []string{
23- `CREATE TABLE IF NOT EXISTS oauth_auth_requests (
24- state TEXT PRIMARY KEY,
25- data TEXT NOT NULL
26- )`,
27- `CREATE TABLE IF NOT EXISTS oauth_sessions (
28- account_did TEXT NOT NULL,
29- session_id TEXT NOT NULL,
30- data TEXT NOT NULL,
31- PRIMARY KEY (account_did, session_id)
32- )`,
33- }
34- for _, stmt := range stmts {
35- if _, err := s.db.ExecContext(ctx, stmt); err != nil {
36- return err
37- }
38- }
39- return nil
40-}
41-
42 func (s *OAuthStore) GetSession(ctx context.Context, did syntax.DID, sessionID string) (*oauth.ClientSessionData, error) {21 func (s *OAuthStore) GetSession(ctx context.Context, did syntax.DID, sessionID string) (*oauth.ClientSessionData, error) {
43 var data []byte22 var data []byte
44 err := s.db.QueryRowContext(ctx, `23 err := s.db.QueryRowContext(ctx, `
@@ -77,6 +56,26 @@ func (s *OAuthStore) DeleteSession(ctx context.Context, did syntax.DID, sessionI
77 return err56 return err
78 }57 }
79 58
59+func (s *OAuthStore) ListSessionsForDID(ctx context.Context, did string) ([]string, error) {
60+ rows, err := s.db.QueryContext(ctx, `
61+ SELECT session_id FROM oauth_sessions WHERE account_did = ? ORDER BY ROWID DESC
62+ `, did)
63+ if err != nil {
64+ return nil, err
65+ }
66+ defer rows.Close()
67+
68+ var ids []string
69+ for rows.Next() {
70+ var id string
71+ if err := rows.Scan(&id); err != nil {
72+ return nil, err
73+ }
74+ ids = append(ids, id)
75+ }
76+ return ids, rows.Err()
77+}
78+
80 func (s *OAuthStore) GetAuthRequestInfo(ctx context.Context, state string) (*oauth.AuthRequestData, error) {79 func (s *OAuthStore) GetAuthRequestInfo(ctx context.Context, state string) (*oauth.AuthRequestData, error) {
81 var data []byte80 var data []byte
82 err := s.db.QueryRowContext(ctx, `81 err := s.db.QueryRowContext(ctx, `
modified internal/db/social.go +12 -0
@@ -152,6 +152,18 @@ func (db *DB) GetLikeCount(ctx context.Context, feedURL, articleURL string) (int
152152 return count, err
153153 }
154154
155+func (db *DB) GetLike(ctx context.Context, authorDID, feedURL, articleURL string) (*Like, error) {
156+ l := &Like{}
157+ err := db.QueryRowContext(ctx, `
158+ SELECT id, uri, author_did, feed_url, article_url, created_at, cid FROM likes
159+ WHERE author_did = ? AND feed_url = ? AND article_url = ?
160+ `, authorDID, feedURL, articleURL).Scan(&l.ID, &l.URI, &l.AuthorDID, &l.FeedURL, &l.ArticleURL, &l.CreatedAt, &l.CID)
161+ if err != nil {
162+ return nil, err
163+ }
164+ return l, nil
165+}
166+
155167 func (db *DB) HasLiked(ctx context.Context, authorDID, feedURL, articleURL string) (bool, error) {
156168 var exists int
157169 err := db.QueryRowContext(ctx, `
@@ -152,6 +152,18 @@ func (db *DB) GetLikeCount(ctx context.Context, feedURL, articleURL string) (int
152 return count, err152 return count, err
153 }153 }
154 154
155+func (db *DB) GetLike(ctx context.Context, authorDID, feedURL, articleURL string) (*Like, error) {
156+ l := &Like{}
157+ err := db.QueryRowContext(ctx, `
158+ SELECT id, uri, author_did, feed_url, article_url, created_at, cid FROM likes
159+ WHERE author_did = ? AND feed_url = ? AND article_url = ?
160+ `, authorDID, feedURL, articleURL).Scan(&l.ID, &l.URI, &l.AuthorDID, &l.FeedURL, &l.ArticleURL, &l.CreatedAt, &l.CID)
161+ if err != nil {
162+ return nil, err
163+ }
164+ return l, nil
165+}
166+
155 func (db *DB) HasLiked(ctx context.Context, authorDID, feedURL, articleURL string) (bool, error) {167 func (db *DB) HasLiked(ctx context.Context, authorDID, feedURL, articleURL string) (bool, error) {
156 var exists int168 var exists int
157 err := db.QueryRowContext(ctx, `169 err := db.QueryRowContext(ctx, `
modified internal/db/user.go +21 -0
@@ -53,3 +53,24 @@ func (db *DB) GetUserByHandle(ctx context.Context, handle string) (*User, error)
5353 }
5454 return u, nil
5555 }
56+
57+func (db *DB) ListUsers(ctx context.Context) ([]*User, error) {
58+ rows, err := db.QueryContext(ctx, `
59+ SELECT did, handle, display_name, avatar_url, indexed_at, updated_at
60+ FROM users ORDER BY updated_at DESC
61+ `)
62+ if err != nil {
63+ return nil, err
64+ }
65+ defer rows.Close()
66+
67+ var users []*User
68+ for rows.Next() {
69+ u := &User{}
70+ if err := rows.Scan(&u.DID, &u.Handle, &u.DisplayName, &u.AvatarURL, &u.IndexedAt, &u.UpdatedAt); err != nil {
71+ return nil, err
72+ }
73+ users = append(users, u)
74+ }
75+ return users, rows.Err()
76+}
@@ -53,3 +53,24 @@ func (db *DB) GetUserByHandle(ctx context.Context, handle string) (*User, error)
53 }53 }
54 return u, nil54 return u, nil
55 }55 }
56+
57+func (db *DB) ListUsers(ctx context.Context) ([]*User, error) {
58+ rows, err := db.QueryContext(ctx, `
59+ SELECT did, handle, display_name, avatar_url, indexed_at, updated_at
60+ FROM users ORDER BY updated_at DESC
61+ `)
62+ if err != nil {
63+ return nil, err
64+ }
65+ defer rows.Close()
66+
67+ var users []*User
68+ for rows.Next() {
69+ u := &User{}
70+ if err := rows.Scan(&u.DID, &u.Handle, &u.DisplayName, &u.AvatarURL, &u.IndexedAt, &u.UpdatedAt); err != nil {
71+ return nil, err
72+ }
73+ users = append(users, u)
74+ }
75+ return users, rows.Err()
76+}
modified internal/server/annotations_handler.go +14 -8
@@ -24,7 +24,6 @@ func (s *Server) handleAnnotations(w http.ResponseWriter, r *http.Request) {
2424 func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request) {
2525 user := currentUser(r)
2626 a := &db.Annotation{
27- URI: fmt.Sprintf("glean:annotation:%d", time.Now().UnixNano()),
2827 AuthorDID: user.DID,
2928 FeedURL: r.FormValue("feed_url"),
3029 ArticleURL: r.FormValue("article_url"),
@@ -40,11 +39,6 @@ func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request)
4039 }
4140 }
4241
43- if err := s.db.CreateAnnotation(r.Context(), a); err != nil {
44- http.Error(w, err.Error(), http.StatusInternalServerError)
45- return
46- }
47-
4842 if client := s.pdsClientForUser(r); client != nil {
4943 record := atproto.AnnotationRecord{
5044 CreatedAt: time.Now().Format(time.RFC3339),
@@ -54,9 +48,21 @@ func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request)
5448 Note: a.Note.String,
5549 Rating: int(a.Rating.Int64),
5650 }
57- if _, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.annotation", record); err != nil {
58- s.logger.Warn("failed to write annotation to PDS", "error", err)
51+ uri, cid, err := client.CreateRecord(r.Context(), user.DID, "at.glean.annotation", record)
52+ if err != nil {
53+ s.logger.Error("failed to write annotation to PDS", "error", err)
54+ http.Error(w, "failed to write annotation to PDS: "+err.Error(), http.StatusBadGateway)
55+ return
5956 }
57+ a.URI = uri
58+ a.CID = sql.NullString{String: cid, Valid: true}
59+ } else {
60+ a.URI = fmt.Sprintf("glean:annotation:%d", time.Now().UnixNano())
61+ }
62+
63+ if err := s.db.CreateAnnotation(r.Context(), a); err != nil {
64+ http.Error(w, err.Error(), http.StatusInternalServerError)
65+ return
6066 }
6167
6268 w.WriteHeader(http.StatusNoContent)
@@ -24,7 +24,6 @@ func (s *Server) handleAnnotations(w http.ResponseWriter, r *http.Request) {
24 func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request) {24 func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request) {
25 user := currentUser(r)25 user := currentUser(r)
26 a := &db.Annotation{26 a := &db.Annotation{
27- URI: fmt.Sprintf("glean:annotation:%d", time.Now().UnixNano()),
28 AuthorDID: user.DID,27 AuthorDID: user.DID,
29 FeedURL: r.FormValue("feed_url"),28 FeedURL: r.FormValue("feed_url"),
30 ArticleURL: r.FormValue("article_url"),29 ArticleURL: r.FormValue("article_url"),
@@ -40,11 +39,6 @@ func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request)
40 }39 }
41 }40 }
42 41
43- if err := s.db.CreateAnnotation(r.Context(), a); err != nil {
44- http.Error(w, err.Error(), http.StatusInternalServerError)
45- return
46- }
47-
48 if client := s.pdsClientForUser(r); client != nil {42 if client := s.pdsClientForUser(r); client != nil {
49 record := atproto.AnnotationRecord{43 record := atproto.AnnotationRecord{
50 CreatedAt: time.Now().Format(time.RFC3339),44 CreatedAt: time.Now().Format(time.RFC3339),
@@ -54,9 +48,21 @@ func (s *Server) handleCreateAnnotation(w http.ResponseWriter, r *http.Request)
54 Note: a.Note.String,48 Note: a.Note.String,
55 Rating: int(a.Rating.Int64),49 Rating: int(a.Rating.Int64),
56 }50 }
57- if _, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.annotation", record); err != nil {51+ uri, cid, err := client.CreateRecord(r.Context(), user.DID, "at.glean.annotation", record)
58- s.logger.Warn("failed to write annotation to PDS", "error", err)52+ if err != nil {
53+ s.logger.Error("failed to write annotation to PDS", "error", err)
54+ http.Error(w, "failed to write annotation to PDS: "+err.Error(), http.StatusBadGateway)
55+ return
59 }56 }
57+ a.URI = uri
58+ a.CID = sql.NullString{String: cid, Valid: true}
59+ } else {
60+ a.URI = fmt.Sprintf("glean:annotation:%d", time.Now().UnixNano())
61+ }
62+
63+ if err := s.db.CreateAnnotation(r.Context(), a); err != nil {
64+ http.Error(w, err.Error(), http.StatusInternalServerError)
65+ return
60 }66 }
61 67
62 w.WriteHeader(http.StatusNoContent)68 w.WriteHeader(http.StatusNoContent)
modified internal/server/articles_handler.go +20 -1
@@ -137,6 +137,23 @@ func (s *Server) handleLikeArticle(w http.ResponseWriter, r *http.Request) {
137137 }
138138
139139 if liked {
140+ existingLike, getErr := s.db.GetLike(r.Context(), user.DID, article.FeedURL, article.URL.String)
141+ if getErr != nil {
142+ http.Error(w, getErr.Error(), http.StatusInternalServerError)
143+ return
144+ }
145+ if existingLike.URI != "" {
146+ if client := s.pdsClientForUser(r); client != nil {
147+ parsed, ok := atproto.ParseRecordURI(existingLike.URI)
148+ if ok {
149+ if delErr := client.DeleteRecord(r.Context(), user.DID, parsed.Collection, parsed.RKey); delErr != nil {
150+ s.logger.Error("failed to delete like from PDS", "error", delErr)
151+ http.Error(w, "failed to delete like from PDS: "+delErr.Error(), http.StatusBadGateway)
152+ return
153+ }
154+ }
155+ }
156+ }
140157 if err := s.db.DeleteLikeByUserArticle(r.Context(), user.DID, article.FeedURL, article.URL.String); err != nil {
141158 http.Error(w, err.Error(), http.StatusInternalServerError)
142159 return
@@ -151,7 +168,9 @@ func (s *Server) handleLikeArticle(w http.ResponseWriter, r *http.Request) {
151168 if client := s.pdsClientForUser(r); client != nil {
152169 uri, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.like", likeRecord)
153170 if err != nil {
154- s.logger.Warn("failed to write like to PDS", "error", err)
171+ s.logger.Error("failed to write like to PDS", "error", err)
172+ http.Error(w, "failed to write like to PDS: "+err.Error(), http.StatusBadGateway)
173+ return
155174 }
156175
157176 like := &db.Like{
@@ -137,6 +137,23 @@ func (s *Server) handleLikeArticle(w http.ResponseWriter, r *http.Request) {
137 }137 }
138 138
139 if liked {139 if liked {
140+ existingLike, getErr := s.db.GetLike(r.Context(), user.DID, article.FeedURL, article.URL.String)
141+ if getErr != nil {
142+ http.Error(w, getErr.Error(), http.StatusInternalServerError)
143+ return
144+ }
145+ if existingLike.URI != "" {
146+ if client := s.pdsClientForUser(r); client != nil {
147+ parsed, ok := atproto.ParseRecordURI(existingLike.URI)
148+ if ok {
149+ if delErr := client.DeleteRecord(r.Context(), user.DID, parsed.Collection, parsed.RKey); delErr != nil {
150+ s.logger.Error("failed to delete like from PDS", "error", delErr)
151+ http.Error(w, "failed to delete like from PDS: "+delErr.Error(), http.StatusBadGateway)
152+ return
153+ }
154+ }
155+ }
156+ }
140 if err := s.db.DeleteLikeByUserArticle(r.Context(), user.DID, article.FeedURL, article.URL.String); err != nil {157 if err := s.db.DeleteLikeByUserArticle(r.Context(), user.DID, article.FeedURL, article.URL.String); err != nil {
141 http.Error(w, err.Error(), http.StatusInternalServerError)158 http.Error(w, err.Error(), http.StatusInternalServerError)
142 return159 return
@@ -151,7 +168,9 @@ func (s *Server) handleLikeArticle(w http.ResponseWriter, r *http.Request) {
151 if client := s.pdsClientForUser(r); client != nil {168 if client := s.pdsClientForUser(r); client != nil {
152 uri, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.like", likeRecord)169 uri, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.like", likeRecord)
153 if err != nil {170 if err != nil {
154- s.logger.Warn("failed to write like to PDS", "error", err)171+ s.logger.Error("failed to write like to PDS", "error", err)
172+ http.Error(w, "failed to write like to PDS: "+err.Error(), http.StatusBadGateway)
173+ return
155 }174 }
156 175
157 like := &db.Like{176 like := &db.Like{
modified internal/server/auth_handler.go +11 -0
@@ -124,6 +124,8 @@ func (s *Server) handleOAuthCallback(w http.ResponseWriter, r *http.Request) {
124124 SameSite: http.SameSiteLaxMode,
125125 })
126126
127+ s.syncUserInBackground(user.DID, s.pdsClientFromSession(sessData))
128+
127129 http.Redirect(w, r, "/dashboard", http.StatusSeeOther)
128130 }
129131
@@ -149,6 +151,15 @@ func (s *Server) fetchUserProfile(ctx context.Context, sessData *oauth.ClientSes
149151 return profile.DisplayName, profile.Avatar
150152 }
151153
154+func (s *Server) pdsClientFromSession(sessData *oauth.ClientSessionData) *atproto.Client {
155+ session, err := s.oauth.ResumeSession(context.Background(), sessData.AccountDID, sessData.SessionID)
156+ if err != nil {
157+ s.logger.Warn("failed to resume session for sync", "error", err)
158+ return nil
159+ }
160+ return &atproto.Client{APIClient: session.APIClient()}
161+}
162+
152163 func (s *Server) handleOAuthClientMetadata(w http.ResponseWriter, r *http.Request) {
153164 if s.clientID == "" {
154165 http.Error(w, "localhost client", http.StatusNotFound)
@@ -124,6 +124,8 @@ func (s *Server) handleOAuthCallback(w http.ResponseWriter, r *http.Request) {
124 SameSite: http.SameSiteLaxMode,124 SameSite: http.SameSiteLaxMode,
125 })125 })
126 126
127+ s.syncUserInBackground(user.DID, s.pdsClientFromSession(sessData))
128+
127 http.Redirect(w, r, "/dashboard", http.StatusSeeOther)129 http.Redirect(w, r, "/dashboard", http.StatusSeeOther)
128 }130 }
129 131
@@ -149,6 +151,15 @@ func (s *Server) fetchUserProfile(ctx context.Context, sessData *oauth.ClientSes
149 return profile.DisplayName, profile.Avatar151 return profile.DisplayName, profile.Avatar
150 }152 }
151 153
154+func (s *Server) pdsClientFromSession(sessData *oauth.ClientSessionData) *atproto.Client {
155+ session, err := s.oauth.ResumeSession(context.Background(), sessData.AccountDID, sessData.SessionID)
156+ if err != nil {
157+ s.logger.Warn("failed to resume session for sync", "error", err)
158+ return nil
159+ }
160+ return &atproto.Client{APIClient: session.APIClient()}
161+}
162+
152 func (s *Server) handleOAuthClientMetadata(w http.ResponseWriter, r *http.Request) {163 func (s *Server) handleOAuthClientMetadata(w http.ResponseWriter, r *http.Request) {
153 if s.clientID == "" {164 if s.clientID == "" {
154 http.Error(w, "localhost client", http.StatusNotFound)165 http.Error(w, "localhost client", http.StatusNotFound)
modified internal/server/feeds_handler.go +45 -13
@@ -85,22 +85,33 @@ func (s *Server) handleAddFeed(w http.ResponseWriter, r *http.Request) {
8585 }
8686 }
8787
88- if err := s.db.CreateSubscription(r.Context(), user.DID, feedURL, category); err != nil {
89- s.logger.Error("failed to create subscription", "error", err)
90- http.Error(w, err.Error(), http.StatusInternalServerError)
91- return
88+ var feedTitle string
89+ if result != nil {
90+ feedTitle = result.Feed.Title
9291 }
9392
93+ var subURI, subCID string
9494 if client := s.pdsClientForUser(r); client != nil {
9595 record := atproto.SubscriptionRecord{
9696 CreatedAt: time.Now().Format(time.RFC3339),
9797 FeedURL: feedURL,
98- Title: result.Feed.Title,
98+ Title: feedTitle,
9999 Category: category,
100100 }
101- if _, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record); err != nil {
102- s.logger.Warn("failed to write subscription to PDS", "error", err)
101+ uri, cid, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record)
102+ if err != nil {
103+ s.logger.Error("failed to write subscription to PDS", "error", err)
104+ http.Error(w, "failed to write subscription to PDS: "+err.Error(), http.StatusBadGateway)
105+ return
103106 }
107+ subURI = uri
108+ subCID = cid
109+ }
110+
111+ if err := s.db.CreateSubscription(r.Context(), user.DID, feedURL, category, subURI, subCID); err != nil {
112+ s.logger.Error("failed to create subscription", "error", err)
113+ http.Error(w, err.Error(), http.StatusInternalServerError)
114+ return
104115 }
105116
106117 subs, _ := s.db.ListSubscriptions(r.Context(), user.DID, "", 100, 0)
@@ -119,6 +130,20 @@ func (s *Server) handleRemoveFeed(w http.ResponseWriter, r *http.Request) {
119130 return
120131 }
121132
133+ sub, err := s.db.GetSubscription(r.Context(), user.DID, feedURL)
134+ if err == nil && sub.URI.Valid {
135+ if client := s.pdsClientForUser(r); client != nil {
136+ parsed, ok := atproto.ParseRecordURI(sub.URI.String)
137+ if ok {
138+ if delErr := client.DeleteRecord(r.Context(), user.DID, parsed.Collection, parsed.RKey); delErr != nil {
139+ s.logger.Error("failed to delete subscription from PDS", "error", delErr)
140+ http.Error(w, "failed to delete subscription from PDS: "+delErr.Error(), http.StatusBadGateway)
141+ return
142+ }
143+ }
144+ }
145+ }
146+
122147 if err := s.db.DeleteSubscription(r.Context(), user.DID, feedURL); err != nil {
123148 s.logger.Error("failed to delete subscription", "error", err)
124149 http.Error(w, err.Error(), http.StatusInternalServerError)
@@ -155,10 +180,8 @@ func (s *Server) handleOPMLUpload(w http.ResponseWriter, r *http.Request) {
155180 s.logger.Error("failed to upsert feed", "error", upsertErr)
156181 continue
157182 }
158- if subErr := s.db.CreateSubscription(r.Context(), user.DID, fu.URL, fu.Category); subErr != nil {
159- s.logger.Error("failed to create subscription", "error", subErr)
160- continue
161- }
183+
184+ var subURI, subCID string
162185 if client != nil {
163186 record := atproto.SubscriptionRecord{
164187 CreatedAt: time.Now().Format(time.RFC3339),
@@ -166,9 +189,18 @@ func (s *Server) handleOPMLUpload(w http.ResponseWriter, r *http.Request) {
166189 Title: fu.Title,
167190 Category: fu.Category,
168191 }
169- if _, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record); err != nil {
170- s.logger.Warn("failed to write subscription to PDS", "error", err, "url", fu.URL)
192+ uri, cid, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record)
193+ if err != nil {
194+ s.logger.Error("failed to write subscription to PDS", "error", err, "url", fu.URL)
195+ continue
171196 }
197+ subURI = uri
198+ subCID = cid
199+ }
200+
201+ if subErr := s.db.CreateSubscription(r.Context(), user.DID, fu.URL, fu.Category, subURI, subCID); subErr != nil {
202+ s.logger.Error("failed to create subscription", "error", subErr)
203+ continue
172204 }
173205 added++
174206 }
@@ -85,22 +85,33 @@ func (s *Server) handleAddFeed(w http.ResponseWriter, r *http.Request) {
85 }85 }
86 }86 }
87 87
88- if err := s.db.CreateSubscription(r.Context(), user.DID, feedURL, category); err != nil {88+ var feedTitle string
89- s.logger.Error("failed to create subscription", "error", err)89+ if result != nil {
90- http.Error(w, err.Error(), http.StatusInternalServerError)90+ feedTitle = result.Feed.Title
91- return
92 }91 }
93 92
93+ var subURI, subCID string
94 if client := s.pdsClientForUser(r); client != nil {94 if client := s.pdsClientForUser(r); client != nil {
95 record := atproto.SubscriptionRecord{95 record := atproto.SubscriptionRecord{
96 CreatedAt: time.Now().Format(time.RFC3339),96 CreatedAt: time.Now().Format(time.RFC3339),
97 FeedURL: feedURL,97 FeedURL: feedURL,
98- Title: result.Feed.Title,98+ Title: feedTitle,
99 Category: category,99 Category: category,
100 }100 }
101- if _, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record); err != nil {101+ uri, cid, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record)
102- s.logger.Warn("failed to write subscription to PDS", "error", err)102+ if err != nil {
103+ s.logger.Error("failed to write subscription to PDS", "error", err)
104+ http.Error(w, "failed to write subscription to PDS: "+err.Error(), http.StatusBadGateway)
105+ return
103 }106 }
107+ subURI = uri
108+ subCID = cid
109+ }
110+
111+ if err := s.db.CreateSubscription(r.Context(), user.DID, feedURL, category, subURI, subCID); err != nil {
112+ s.logger.Error("failed to create subscription", "error", err)
113+ http.Error(w, err.Error(), http.StatusInternalServerError)
114+ return
104 }115 }
105 116
106 subs, _ := s.db.ListSubscriptions(r.Context(), user.DID, "", 100, 0)117 subs, _ := s.db.ListSubscriptions(r.Context(), user.DID, "", 100, 0)
@@ -119,6 +130,20 @@ func (s *Server) handleRemoveFeed(w http.ResponseWriter, r *http.Request) {
119 return130 return
120 }131 }
121 132
133+ sub, err := s.db.GetSubscription(r.Context(), user.DID, feedURL)
134+ if err == nil && sub.URI.Valid {
135+ if client := s.pdsClientForUser(r); client != nil {
136+ parsed, ok := atproto.ParseRecordURI(sub.URI.String)
137+ if ok {
138+ if delErr := client.DeleteRecord(r.Context(), user.DID, parsed.Collection, parsed.RKey); delErr != nil {
139+ s.logger.Error("failed to delete subscription from PDS", "error", delErr)
140+ http.Error(w, "failed to delete subscription from PDS: "+delErr.Error(), http.StatusBadGateway)
141+ return
142+ }
143+ }
144+ }
145+ }
146+
122 if err := s.db.DeleteSubscription(r.Context(), user.DID, feedURL); err != nil {147 if err := s.db.DeleteSubscription(r.Context(), user.DID, feedURL); err != nil {
123 s.logger.Error("failed to delete subscription", "error", err)148 s.logger.Error("failed to delete subscription", "error", err)
124 http.Error(w, err.Error(), http.StatusInternalServerError)149 http.Error(w, err.Error(), http.StatusInternalServerError)
@@ -155,10 +180,8 @@ func (s *Server) handleOPMLUpload(w http.ResponseWriter, r *http.Request) {
155 s.logger.Error("failed to upsert feed", "error", upsertErr)180 s.logger.Error("failed to upsert feed", "error", upsertErr)
156 continue181 continue
157 }182 }
158- if subErr := s.db.CreateSubscription(r.Context(), user.DID, fu.URL, fu.Category); subErr != nil {183+
159- s.logger.Error("failed to create subscription", "error", subErr)184+ var subURI, subCID string
160- continue
161- }
162 if client != nil {185 if client != nil {
163 record := atproto.SubscriptionRecord{186 record := atproto.SubscriptionRecord{
164 CreatedAt: time.Now().Format(time.RFC3339),187 CreatedAt: time.Now().Format(time.RFC3339),
@@ -166,9 +189,18 @@ func (s *Server) handleOPMLUpload(w http.ResponseWriter, r *http.Request) {
166 Title: fu.Title,189 Title: fu.Title,
167 Category: fu.Category,190 Category: fu.Category,
168 }191 }
169- if _, _, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record); err != nil {192+ uri, cid, err := client.CreateRecord(r.Context(), user.DID, "at.glean.subscription", record)
170- s.logger.Warn("failed to write subscription to PDS", "error", err, "url", fu.URL)193+ if err != nil {
194+ s.logger.Error("failed to write subscription to PDS", "error", err, "url", fu.URL)
195+ continue
171 }196 }
197+ subURI = uri
198+ subCID = cid
199+ }
200+
201+ if subErr := s.db.CreateSubscription(r.Context(), user.DID, fu.URL, fu.Category, subURI, subCID); subErr != nil {
202+ s.logger.Error("failed to create subscription", "error", subErr)
203+ continue
172 }204 }
173 added++205 added++
174 }206 }
modified internal/server/server.go +65 -11
@@ -29,22 +29,19 @@ func splitString(s, sep string) []string {
2929 }
3030
3131 type Server struct {
32- db *db.DB
33- router *chi.Mux
34- templates *template.Template
35- logger *slog.Logger
36- oauth *oauth.ClientApp
37- oauthStore *db.OAuthStore
38- fetcher *feed.Fetcher
39- clientID string
32+ db *db.DB
33+ router *chi.Mux
34+ templates *template.Template
35+ logger *slog.Logger
36+ oauth *oauth.ClientApp
37+ oauthStore *db.OAuthStore
38+ fetcher *feed.Fetcher
39+ clientID string
4040 callbackURL string
4141 }
4242
4343 func New(database *db.DB, clientID, callbackURL, addr string, logger *slog.Logger) *Server {
4444 oauthStore := db.NewOAuthStore(database)
45- if err := oauthStore.Init(context.Background()); err != nil {
46- logger.Error("failed to init oauth store", "error", err)
47- }
4845
4946 var config oauth.ClientConfig
5047 if clientID == "" {
@@ -258,6 +255,63 @@ func (s *Server) pdsClientForUser(r *http.Request) *atproto.Client {
258255 return nil
259256 }
260257
258+func (s *Server) syncUserInBackground(userDID string, client *atproto.Client) {
259+ go func() {
260+ ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
261+ defer cancel()
262+ sync := atproto.NewSync(s.db, client, s.logger)
263+ if err := sync.Run(ctx, userDID); err != nil {
264+ s.logger.Error("background sync failed", "error", err, "did", userDID)
265+ }
266+ }()
267+}
268+
269+func (s *Server) PeriodicSync(ctx context.Context, interval time.Duration) {
270+ ticker := time.NewTicker(interval)
271+ defer ticker.Stop()
272+
273+ for {
274+ select {
275+ case <-ctx.Done():
276+ return
277+ case <-ticker.C:
278+ s.runSyncAll(ctx)
279+ }
280+ }
281+}
282+
283+func (s *Server) runSyncAll(ctx context.Context) {
284+ users, err := s.db.ListUsers(ctx)
285+ if err != nil {
286+ s.logger.Error("failed to list users for sync", "error", err)
287+ return
288+ }
289+
290+ for _, u := range users {
291+ sessionIDs, err := s.oauthStore.ListSessionsForDID(ctx, u.DID)
292+ if err != nil || len(sessionIDs) == 0 {
293+ continue
294+ }
295+
296+ did, err := syntax.ParseDID(u.DID)
297+ if err != nil {
298+ continue
299+ }
300+
301+ sess, err := s.oauth.ResumeSession(ctx, did, sessionIDs[0])
302+ if err != nil {
303+ s.logger.Warn("failed to resume session for periodic sync", "error", err, "did", u.DID)
304+ continue
305+ }
306+
307+ client := &atproto.Client{APIClient: sess.APIClient()}
308+ sync := atproto.NewSync(s.db, client, s.logger)
309+ if err := sync.Run(ctx, u.DID); err != nil {
310+ s.logger.Error("periodic sync failed", "error", err, "did", u.DID)
311+ }
312+ }
313+}
314+
261315 func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
262316 s.router.ServeHTTP(w, r)
263317 }
@@ -29,22 +29,19 @@ func splitString(s, sep string) []string {
29 }29 }
30 30
31 type Server struct {31 type Server struct {
32- db *db.DB32+ db *db.DB
33- router *chi.Mux33+ router *chi.Mux
34- templates *template.Template34+ templates *template.Template
35- logger *slog.Logger35+ logger *slog.Logger
36- oauth *oauth.ClientApp36+ oauth *oauth.ClientApp
37- oauthStore *db.OAuthStore37+ oauthStore *db.OAuthStore
38- fetcher *feed.Fetcher38+ fetcher *feed.Fetcher
39- clientID string39+ clientID string
40 callbackURL string40 callbackURL string
41 }41 }
42 42
43 func New(database *db.DB, clientID, callbackURL, addr string, logger *slog.Logger) *Server {43 func New(database *db.DB, clientID, callbackURL, addr string, logger *slog.Logger) *Server {
44 oauthStore := db.NewOAuthStore(database)44 oauthStore := db.NewOAuthStore(database)
45- if err := oauthStore.Init(context.Background()); err != nil {
46- logger.Error("failed to init oauth store", "error", err)
47- }
48 45
49 var config oauth.ClientConfig46 var config oauth.ClientConfig
50 if clientID == "" {47 if clientID == "" {
@@ -258,6 +255,63 @@ func (s *Server) pdsClientForUser(r *http.Request) *atproto.Client {
258 return nil255 return nil
259 }256 }
260 257
258+func (s *Server) syncUserInBackground(userDID string, client *atproto.Client) {
259+ go func() {
260+ ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
261+ defer cancel()
262+ sync := atproto.NewSync(s.db, client, s.logger)
263+ if err := sync.Run(ctx, userDID); err != nil {
264+ s.logger.Error("background sync failed", "error", err, "did", userDID)
265+ }
266+ }()
267+}
268+
269+func (s *Server) PeriodicSync(ctx context.Context, interval time.Duration) {
270+ ticker := time.NewTicker(interval)
271+ defer ticker.Stop()
272+
273+ for {
274+ select {
275+ case <-ctx.Done():
276+ return
277+ case <-ticker.C:
278+ s.runSyncAll(ctx)
279+ }
280+ }
281+}
282+
283+func (s *Server) runSyncAll(ctx context.Context) {
284+ users, err := s.db.ListUsers(ctx)
285+ if err != nil {
286+ s.logger.Error("failed to list users for sync", "error", err)
287+ return
288+ }
289+
290+ for _, u := range users {
291+ sessionIDs, err := s.oauthStore.ListSessionsForDID(ctx, u.DID)
292+ if err != nil || len(sessionIDs) == 0 {
293+ continue
294+ }
295+
296+ did, err := syntax.ParseDID(u.DID)
297+ if err != nil {
298+ continue
299+ }
300+
301+ sess, err := s.oauth.ResumeSession(ctx, did, sessionIDs[0])
302+ if err != nil {
303+ s.logger.Warn("failed to resume session for periodic sync", "error", err, "did", u.DID)
304+ continue
305+ }
306+
307+ client := &atproto.Client{APIClient: sess.APIClient()}
308+ sync := atproto.NewSync(s.db, client, s.logger)
309+ if err := sync.Run(ctx, u.DID); err != nil {
310+ s.logger.Error("periodic sync failed", "error", err, "did", u.DID)
311+ }
312+ }
313+}
314+
261 func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {315 func (s *Server) ServeHTTP(w http.ResponseWriter, r *http.Request) {
262 s.router.ServeHTTP(w, r)316 s.router.ServeHTTP(w, r)
263 }317 }
modified internal/tmpl/base.html +6 -15
@@ -18,18 +18,9 @@
1818 <body class="bg-spot-bg text-spot-text min-h-screen flex">
1919 {{if .User}}
2020 <aside class="hidden lg:flex flex-col w-60 bg-spot-bg h-screen fixed left-0 top-0 px-3 py-4 z-20">
21- <a href="/" class="text-spot-green font-bold text-xl tracking-tight mb-8 px-3 flex items-center gap-2">
22- <svg class="w-7 h-7" viewBox="0 0 32 32" fill="none" xmlns="http://www.w3.org/2000/svg">
23- <rect width="32" height="32" rx="8" fill="#00754A"/>
24- <path d="M16 8 L16 22" stroke="#fff" stroke-width="2.5" stroke-linecap="round"/>
25- <path d="M16 11 Q11 9 9 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
26- <path d="M16 11 Q21 9 23 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
27- <path d="M16 15 Q10 13 8 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
28- <path d="M16 15 Q22 13 24 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
29- <circle cx="16" cy="24" r="1.5" fill="#fff"/>
30- </svg>
31- Glean
32- </a>
21+ <div class="mb-8 px-3">
22+ {{template "logo-link"}}
23+ </div>
3324 <nav class="flex flex-col gap-1 text-sm">
3425 <a href="/dashboard" class="sidebar-link flex items-center gap-3 px-3 py-2 rounded-md {{activeClass .ActivePath "/dashboard"}}">
3526 <svg class="w-5 h-5" fill="none" stroke="currentColor" viewBox="0 0 24 24"><path stroke-linecap="round" stroke-linejoin="round" stroke-width="2" d="M3 12l2-2m0 0l7-7 7 7M5 10v10a1 1 0 001 1h3m10-11l2 2m-2-2v10a1 1 0 01-1 1h-3m-4 0a1 1 0 01-1-1v-4a1 1 0 011-1h2a1 1 0 011 1v4a1 1 0 01-1 1"/></svg>
@@ -101,7 +92,7 @@
10192 <main id="main-content" class="{{if .User}}lg:ml-60 pb-20 lg:pb-0{{end}} flex-1 min-h-screen flex flex-col">
10293 {{if .User}}
10394 <div class="lg:hidden bg-spot-surface border-b border-spot-divider px-4 py-3 flex items-center justify-between sticky top-0 z-20">
104- <a href="/" class="text-spot-green font-bold text-lg font-title">Glean</a>
95+ {{template "logo-text"}}
10596 <a href="/profile/{{.User.DID}}" class="flex items-center gap-2">
10697 {{if .User.AvatarURL.Valid}}<img src="{{.User.AvatarURL.String}}" class="w-7 h-7 rounded-full">{{end}}
10798 <span class="text-xs text-spot-secondary">@{{.User.Handle}}</span>
@@ -116,7 +107,7 @@
116107 <div class="max-w-6xl mx-auto px-4 lg:px-8 py-8">
117108 <div class="hidden md:grid md:grid-cols-[auto_1fr_auto] md:gap-12 md:items-start">
118109 <div>
119- <a href="/" class="text-spot-green font-bold text-lg tracking-tight">Glean</a>
110+ {{template "logo-text"}}
120111 <p class="text-xs text-spot-secondary mt-1 max-w-[200px]">A social RSS reader<br>on the AT Protocol.</p>
121112 </div>
122113 <div class="grid grid-cols-3 gap-8">
@@ -167,7 +158,7 @@
167158 </div>
168159
169160 <div class="md:hidden flex flex-col gap-6">
170- <a href="/" class="text-spot-green font-bold text-lg tracking-tight">Glean</a>
161+ {{template "logo-text"}}
171162 <div class="grid grid-cols-2 gap-6">
172163 <div class="flex flex-col gap-1.5">
173164 <div class="text-spot-text font-bold text-xs uppercase tracking-wide mb-1">Browse</div>
@@ -18,18 +18,9 @@
18 <body class="bg-spot-bg text-spot-text min-h-screen flex">18 <body class="bg-spot-bg text-spot-text min-h-screen flex">
19 {{if .User}}19 {{if .User}}
20 <aside class="hidden lg:flex flex-col w-60 bg-spot-bg h-screen fixed left-0 top-0 px-3 py-4 z-20">20 <aside class="hidden lg:flex flex-col w-60 bg-spot-bg h-screen fixed left-0 top-0 px-3 py-4 z-20">
21- <a href="/" class="text-spot-green font-bold text-xl tracking-tight mb-8 px-3 flex items-center gap-2">21+ <div class="mb-8 px-3">
22- <svg class="w-7 h-7" viewBox="0 0 32 32" fill="none" xmlns="http://www.w3.org/2000/svg">22+ {{template "logo-link"}}
23- <rect width="32" height="32" rx="8" fill="#00754A"/>23+ </div>
24- <path d="M16 8 L16 22" stroke="#fff" stroke-width="2.5" stroke-linecap="round"/>
25- <path d="M16 11 Q11 9 9 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
26- <path d="M16 11 Q21 9 23 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
27- <path d="M16 15 Q10 13 8 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
28- <path d="M16 15 Q22 13 24 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
29- <circle cx="16" cy="24" r="1.5" fill="#fff"/>
30- </svg>
31- Glean
32- </a>
33 <nav class="flex flex-col gap-1 text-sm">24 <nav class="flex flex-col gap-1 text-sm">
34 <a href="/dashboard" class="sidebar-link flex items-center gap-3 px-3 py-2 rounded-md {{activeClass .ActivePath "/dashboard"}}">25 <a href="/dashboard" class="sidebar-link flex items-center gap-3 px-3 py-2 rounded-md {{activeClass .ActivePath "/dashboard"}}">
35 <svg class="w-5 h-5" fill="none" stroke="currentColor" viewBox="0 0 24 24"><path stroke-linecap="round" stroke-linejoin="round" stroke-width="2" d="M3 12l2-2m0 0l7-7 7 7M5 10v10a1 1 0 001 1h3m10-11l2 2m-2-2v10a1 1 0 01-1 1h-3m-4 0a1 1 0 01-1-1v-4a1 1 0 011-1h2a1 1 0 011 1v4a1 1 0 01-1 1"/></svg>26 <svg class="w-5 h-5" fill="none" stroke="currentColor" viewBox="0 0 24 24"><path stroke-linecap="round" stroke-linejoin="round" stroke-width="2" d="M3 12l2-2m0 0l7-7 7 7M5 10v10a1 1 0 001 1h3m10-11l2 2m-2-2v10a1 1 0 01-1 1h-3m-4 0a1 1 0 01-1-1v-4a1 1 0 011-1h2a1 1 0 011 1v4a1 1 0 01-1 1"/></svg>
@@ -101,7 +92,7 @@
101 <main id="main-content" class="{{if .User}}lg:ml-60 pb-20 lg:pb-0{{end}} flex-1 min-h-screen flex flex-col">92 <main id="main-content" class="{{if .User}}lg:ml-60 pb-20 lg:pb-0{{end}} flex-1 min-h-screen flex flex-col">
102 {{if .User}}93 {{if .User}}
103 <div class="lg:hidden bg-spot-surface border-b border-spot-divider px-4 py-3 flex items-center justify-between sticky top-0 z-20">94 <div class="lg:hidden bg-spot-surface border-b border-spot-divider px-4 py-3 flex items-center justify-between sticky top-0 z-20">
104- <a href="/" class="text-spot-green font-bold text-lg font-title">Glean</a>95+ {{template "logo-text"}}
105 <a href="/profile/{{.User.DID}}" class="flex items-center gap-2">96 <a href="/profile/{{.User.DID}}" class="flex items-center gap-2">
106 {{if .User.AvatarURL.Valid}}<img src="{{.User.AvatarURL.String}}" class="w-7 h-7 rounded-full">{{end}}97 {{if .User.AvatarURL.Valid}}<img src="{{.User.AvatarURL.String}}" class="w-7 h-7 rounded-full">{{end}}
107 <span class="text-xs text-spot-secondary">@{{.User.Handle}}</span>98 <span class="text-xs text-spot-secondary">@{{.User.Handle}}</span>
@@ -116,7 +107,7 @@
116 <div class="max-w-6xl mx-auto px-4 lg:px-8 py-8">107 <div class="max-w-6xl mx-auto px-4 lg:px-8 py-8">
117 <div class="hidden md:grid md:grid-cols-[auto_1fr_auto] md:gap-12 md:items-start">108 <div class="hidden md:grid md:grid-cols-[auto_1fr_auto] md:gap-12 md:items-start">
118 <div>109 <div>
119- <a href="/" class="text-spot-green font-bold text-lg tracking-tight">Glean</a>110+ {{template "logo-text"}}
120 <p class="text-xs text-spot-secondary mt-1 max-w-[200px]">A social RSS reader<br>on the AT Protocol.</p>111 <p class="text-xs text-spot-secondary mt-1 max-w-[200px]">A social RSS reader<br>on the AT Protocol.</p>
121 </div>112 </div>
122 <div class="grid grid-cols-3 gap-8">113 <div class="grid grid-cols-3 gap-8">
@@ -167,7 +158,7 @@
167 </div>158 </div>
168 159
169 <div class="md:hidden flex flex-col gap-6">160 <div class="md:hidden flex flex-col gap-6">
170- <a href="/" class="text-spot-green font-bold text-lg tracking-tight">Glean</a>161+ {{template "logo-text"}}
171 <div class="grid grid-cols-2 gap-6">162 <div class="grid grid-cols-2 gap-6">
172 <div class="flex flex-col gap-1.5">163 <div class="flex flex-col gap-1.5">
173 <div class="text-spot-text font-bold text-xs uppercase tracking-wide mb-1">Browse</div>164 <div class="text-spot-text font-bold text-xs uppercase tracking-wide mb-1">Browse</div>
modified internal/tmpl/feeds.html +2 -1
@@ -44,7 +44,8 @@
4444 {{range .Subscriptions}}
4545 <div class="px-5 py-4 flex items-center justify-between hover:bg-spot-hover-50 transition rounded-xl">
4646 <div class="min-w-0">
47- <div class="font-bold text-spot-text">{{if .FeedTitle}}{{.FeedTitle}}{{else}}{{.FeedURL}}{{end}}</div>
47+ <div class="font-bold text-spot-text truncate">{{if .FeedTitle}}{{.FeedTitle}}{{else}}{{.FeedURL}}{{end}}</div>
48+ {{if and .FeedTitle .FeedURL}}<div class="text-xs text-spot-muted truncate">{{.FeedURL}}</div>{{end}}
4849 {{if .Category.Valid}}<span class="text-xs text-spot-secondary">{{.Category.String}}</span>{{end}}
4950 </div>
5051 <div class="flex items-center gap-3">
@@ -44,7 +44,8 @@
44 {{range .Subscriptions}}44 {{range .Subscriptions}}
45 <div class="px-5 py-4 flex items-center justify-between hover:bg-spot-hover-50 transition rounded-xl">45 <div class="px-5 py-4 flex items-center justify-between hover:bg-spot-hover-50 transition rounded-xl">
46 <div class="min-w-0">46 <div class="min-w-0">
47- <div class="font-bold text-spot-text">{{if .FeedTitle}}{{.FeedTitle}}{{else}}{{.FeedURL}}{{end}}</div>47+ <div class="font-bold text-spot-text truncate">{{if .FeedTitle}}{{.FeedTitle}}{{else}}{{.FeedURL}}{{end}}</div>
48+ {{if and .FeedTitle .FeedURL}}<div class="text-xs text-spot-muted truncate">{{.FeedURL}}</div>{{end}}
48 {{if .Category.Valid}}<span class="text-xs text-spot-secondary">{{.Category.String}}</span>{{end}}49 {{if .Category.Valid}}<span class="text-xs text-spot-secondary">{{.Category.String}}</span>{{end}}
49 </div>50 </div>
50 <div class="flex items-center gap-3">51 <div class="flex items-center gap-3">
modified internal/tmpl/index.html +1 -9
@@ -9,15 +9,7 @@
99 <div class="grid grid-cols-1 lg:grid-cols-2 gap-12 lg:gap-16 items-center">
1010 <div>
1111 <div class="flex items-center gap-3 mb-8">
12- <svg class="w-10 h-10" viewBox="0 0 32 32" fill="none" xmlns="http://www.w3.org/2000/svg">
13- <rect width="32" height="32" rx="8" fill="#00754A"/>
14- <path d="M16 8 L16 22" stroke="#fff" stroke-width="2.5" stroke-linecap="round"/>
15- <path d="M16 11 Q11 9 9 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
16- <path d="M16 11 Q21 9 23 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
17- <path d="M16 15 Q10 13 8 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
18- <path d="M16 15 Q22 13 24 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
19- <circle cx="16" cy="24" r="1.5" fill="#fff"/>
20- </svg>
12+ <span class="w-10 h-10">{{template "logo-icon"}}</span>
2113 <span class="font-bold text-xl text-spot-text" style="letter-spacing: -0.02em;">Glean</span>
2214 </div>
2315 <h1 class="text-3xl md:text-5xl font-bold mb-6 leading-[1.1] text-spot-text" style="letter-spacing: -0.03em;">
@@ -9,15 +9,7 @@
9 <div class="grid grid-cols-1 lg:grid-cols-2 gap-12 lg:gap-16 items-center">9 <div class="grid grid-cols-1 lg:grid-cols-2 gap-12 lg:gap-16 items-center">
10 <div>10 <div>
11 <div class="flex items-center gap-3 mb-8">11 <div class="flex items-center gap-3 mb-8">
12- <svg class="w-10 h-10" viewBox="0 0 32 32" fill="none" xmlns="http://www.w3.org/2000/svg">12+ <span class="w-10 h-10">{{template "logo-icon"}}</span>
13- <rect width="32" height="32" rx="8" fill="#00754A"/>
14- <path d="M16 8 L16 22" stroke="#fff" stroke-width="2.5" stroke-linecap="round"/>
15- <path d="M16 11 Q11 9 9 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
16- <path d="M16 11 Q21 9 23 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
17- <path d="M16 15 Q10 13 8 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
18- <path d="M16 15 Q22 13 24 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
19- <circle cx="16" cy="24" r="1.5" fill="#fff"/>
20- </svg>
21 <span class="font-bold text-xl text-spot-text" style="letter-spacing: -0.02em;">Glean</span>13 <span class="font-bold text-xl text-spot-text" style="letter-spacing: -0.02em;">Glean</span>
22 </div>14 </div>
23 <h1 class="text-3xl md:text-5xl font-bold mb-6 leading-[1.1] text-spot-text" style="letter-spacing: -0.03em;">15 <h1 class="text-3xl md:text-5xl font-bold mb-6 leading-[1.1] text-spot-text" style="letter-spacing: -0.03em;">
modified internal/tmpl/login.html +5 -2
@@ -1,8 +1,11 @@
11 {{define "login.html"}}
22 <div class="max-w-md mx-auto mt-20">
33 <div class="bg-spot-surface rounded-xl shadow-spot-heavy p-8">
4- <h1 class="text-2xl font-bold text-spot-text mb-2">Sign in to Glean</h1>
5- <p class="text-spot-secondary text-sm mb-6">Enter your handle.</p>
4+ <div class="flex justify-center mb-6">
5+ <span class="w-14 h-14">{{template "logo-icon"}}</span>
6+ </div>
7+ <h1 class="text-2xl font-bold text-spot-text mb-2 text-center">Sign in to Glean</h1>
8+ <p class="text-spot-secondary text-sm mb-6 text-center">Enter your handle.</p>
69
710 <form action="/auth/start" method="POST" id="login-form">
811 {{csrfInput .CSRFToken}}
@@ -1,8 +1,11 @@
1 {{define "login.html"}}1 {{define "login.html"}}
2 <div class="max-w-md mx-auto mt-20">2 <div class="max-w-md mx-auto mt-20">
3 <div class="bg-spot-surface rounded-xl shadow-spot-heavy p-8">3 <div class="bg-spot-surface rounded-xl shadow-spot-heavy p-8">
4- <h1 class="text-2xl font-bold text-spot-text mb-2">Sign in to Glean</h1>4+ <div class="flex justify-center mb-6">
5- <p class="text-spot-secondary text-sm mb-6">Enter your handle.</p>5+ <span class="w-14 h-14">{{template "logo-icon"}}</span>
6+ </div>
7+ <h1 class="text-2xl font-bold text-spot-text mb-2 text-center">Sign in to Glean</h1>
8+ <p class="text-spot-secondary text-sm mb-6 text-center">Enter your handle.</p>
6 9
7 <form action="/auth/start" method="POST" id="login-form">10 <form action="/auth/start" method="POST" id="login-form">
8 {{csrfInput .CSRFToken}}11 {{csrfInput .CSRFToken}}
added internal/tmpl/partials/logo.html +16 -0
new file mode 100644
@@ -0,0 +1,16 @@
1+{{define "logo-icon"}}<svg class="w-full h-full" viewBox="0 0 32 32" fill="none" xmlns="http://www.w3.org/2000/svg">
2+ <rect width="32" height="32" rx="8" fill="#00754A"/>
3+ <path d="M16 8 L16 22" stroke="#fff" stroke-width="2.5" stroke-linecap="round"/>
4+ <path d="M16 11 Q11 9 9 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
5+ <path d="M16 11 Q21 9 23 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
6+ <path d="M16 15 Q10 13 8 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
7+ <path d="M16 15 Q22 13 24 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
8+ <circle cx="16" cy="24" r="1.5" fill="#fff"/>
9+</svg>{{end}}
10+
11+{{define "logo-link"}}<a href="/" class="text-spot-green font-bold text-xl tracking-tight flex items-center gap-2">
12+ <span class="w-7 h-7">{{template "logo-icon"}}</span>
13+ Glean
14+</a>{{end}}
15+
16+{{define "logo-text"}}<a href="/" class="text-spot-green font-bold text-lg tracking-tight">Glean</a>{{end}}
new file mode 100644
@@ -0,0 +1,16 @@
1+{{define "logo-icon"}}<svg class="w-full h-full" viewBox="0 0 32 32" fill="none" xmlns="http://www.w3.org/2000/svg">
2+ <rect width="32" height="32" rx="8" fill="#00754A"/>
3+ <path d="M16 8 L16 22" stroke="#fff" stroke-width="2.5" stroke-linecap="round"/>
4+ <path d="M16 11 Q11 9 9 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
5+ <path d="M16 11 Q21 9 23 13" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
6+ <path d="M16 15 Q10 13 8 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
7+ <path d="M16 15 Q22 13 24 18" stroke="#fff" stroke-width="1.8" stroke-linecap="round" fill="none"/>
8+ <circle cx="16" cy="24" r="1.5" fill="#fff"/>
9+</svg>{{end}}
10+
11+{{define "logo-link"}}<a href="/" class="text-spot-green font-bold text-xl tracking-tight flex items-center gap-2">
12+ <span class="w-7 h-7">{{template "logo-icon"}}</span>
13+ Glean
14+</a>{{end}}
15+
16+{{define "logo-text"}}<a href="/" class="text-spot-green font-bold text-lg tracking-tight">Glean</a>{{end}}
modified main.go +5 -4
@@ -44,10 +44,8 @@ func main() {
4444 engine := cluster.NewEngine(database.DB, logger)
4545 cron := cluster.NewCron(engine, 6*time.Hour, logger)
4646
47- firehose := atproto.NewFirehoseConsumer(*relayURL, func(ctx context.Context, event *atproto.FirehoseEvent) error {
48- logger.Debug("firehose event", "type", event.Type, "collection", event.Collection, "did", event.DID)
49- return nil
50- }, logger)
47+ firehoseHandler := atproto.NewFirehoseDBHandler(database, logger)
48+ firehose := atproto.NewFirehoseConsumer(*relayURL, firehoseHandler.Handle, logger)
5149
5250 ctx, cancel := context.WithCancel(context.Background())
5351 defer cancel()
@@ -62,6 +60,9 @@ func main() {
6260 logger.Error("cron error", "error", err)
6361 }
6462 }()
63+ go func() {
64+ srv.PeriodicSync(ctx, 1*time.Hour)
65+ }()
6566 go func() {
6667 if err := firehose.Start(ctx); err != nil && ctx.Err() == nil {
6768 logger.Error("firehose error", "error", err)
@@ -44,10 +44,8 @@ func main() {
44 engine := cluster.NewEngine(database.DB, logger)44 engine := cluster.NewEngine(database.DB, logger)
45 cron := cluster.NewCron(engine, 6*time.Hour, logger)45 cron := cluster.NewCron(engine, 6*time.Hour, logger)
46 46
47- firehose := atproto.NewFirehoseConsumer(*relayURL, func(ctx context.Context, event *atproto.FirehoseEvent) error {47+ firehoseHandler := atproto.NewFirehoseDBHandler(database, logger)
48- logger.Debug("firehose event", "type", event.Type, "collection", event.Collection, "did", event.DID)48+ firehose := atproto.NewFirehoseConsumer(*relayURL, firehoseHandler.Handle, logger)
49- return nil
50- }, logger)
51 49
52 ctx, cancel := context.WithCancel(context.Background())50 ctx, cancel := context.WithCancel(context.Background())
53 defer cancel()51 defer cancel()
@@ -62,6 +60,9 @@ func main() {
62 logger.Error("cron error", "error", err)60 logger.Error("cron error", "error", err)
63 }61 }
64 }()62 }()
63+ go func() {
64+ srv.PeriodicSync(ctx, 1*time.Hour)
65+ }()
65 go func() {66 go func() {
66 if err := firehose.Start(ctx); err != nil && ctx.Err() == nil {67 if err := firehose.Start(ctx); err != nil && ctx.Err() == nil {
67 logger.Error("firehose error", "error", err)68 logger.Error("firehose error", "error", err)