nandi/gleanpublic Fork 0
041a564
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.

Introduce persistent cursor storage for Jetstream consumerUnverified

Julien Robert committed 2026-05-07T20:50:51+02:00 Browse files
041a564 parent: 71c0a01
modified go.mod +1 -1
@@ -13,7 +13,6 @@ require (
1313 github.com/prometheus/client_golang v1.19.1
1414 github.com/prometheus/client_model v0.6.1
1515 github.com/prometheus/common v0.54.0
16- go.uber.org/atomic v1.11.0
1716 golang.org/x/net v0.53.0
1817 golang.org/x/sync v0.20.0
1918 gotest.tools/v3 v3.5.2
@@ -48,6 +47,7 @@ require (
4847 github.com/whyrusleeping/cbor-gen v0.2.1-0.20241030202151-b7a6831be65e // indirect
4948 gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect
5049 gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect
50+ go.uber.org/atomic v1.11.0 // indirect
5151 golang.org/x/crypto v0.50.0 // indirect
5252 golang.org/x/sys v0.43.0 // indirect
5353 golang.org/x/text v0.36.0 // indirect
@@ -13,7 +13,6 @@ require (
13 github.com/prometheus/client_golang v1.19.113 github.com/prometheus/client_golang v1.19.1
14 github.com/prometheus/client_model v0.6.114 github.com/prometheus/client_model v0.6.1
15 github.com/prometheus/common v0.54.015 github.com/prometheus/common v0.54.0
16- go.uber.org/atomic v1.11.0
17 golang.org/x/net v0.53.016 golang.org/x/net v0.53.0
18 golang.org/x/sync v0.20.017 golang.org/x/sync v0.20.0
19 gotest.tools/v3 v3.5.218 gotest.tools/v3 v3.5.2
@@ -48,6 +47,7 @@ require (
48 github.com/whyrusleeping/cbor-gen v0.2.1-0.20241030202151-b7a6831be65e // indirect47 github.com/whyrusleeping/cbor-gen v0.2.1-0.20241030202151-b7a6831be65e // indirect
49 gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect48 gitlab.com/yawning/secp256k1-voi v0.0.0-20230925100816-f2616030848b // indirect
50 gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect49 gitlab.com/yawning/tuplehash v0.0.0-20230713102510-df83abbf9a02 // indirect
50+ go.uber.org/atomic v1.11.0 // indirect
51 golang.org/x/crypto v0.50.0 // indirect51 golang.org/x/crypto v0.50.0 // indirect
52 golang.org/x/sys v0.43.0 // indirect52 golang.org/x/sys v0.43.0 // indirect
53 golang.org/x/text v0.36.0 // indirect53 golang.org/x/text v0.36.0 // indirect
modified internal/atproto/jetstream.go +40 -18
@@ -10,11 +10,16 @@ import (
1010
1111 jsc "github.com/bluesky-social/jetstream/pkg/client"
1212 "github.com/bluesky-social/jetstream/pkg/models"
13- "go.uber.org/atomic"
1413
1514 "pkg.rbrt.fr/glean/internal/metrics"
1615 )
1716
17+// CursorStore stores the Jetstream cursor.
18+type CursorStore interface {
19+ LoadCursor(ctx context.Context) (*int64, error)
20+ SaveCursor(ctx context.Context, cursor int64) error
21+}
22+
1823 type Event struct {
1924 Type string
2025 DID string
@@ -28,14 +33,16 @@ type Event struct {
2833 type EventHandler func(ctx context.Context, event *Event) error
2934
3035 type jetstreamScheduler struct {
31- handler EventHandler
32- logger *slog.Logger
33- cursor atomic.Int64
36+ handler EventHandler
37+ logger *slog.Logger
38+ cursorStore CursorStore
3439 }
3540
3641 func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.Event) error {
3742 if evt.TimeUS > 0 {
38- s.cursor.Store(evt.TimeUS)
43+ if err := s.cursorStore.SaveCursor(ctx, evt.TimeUS); err != nil {
44+ s.logger.Warn("failed to save cursor", "error", err)
45+ }
3946 }
4047
4148 if evt.Kind != models.EventKindCommit || evt.Commit == nil {
@@ -71,15 +78,18 @@ func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.
7178 func (s *jetstreamScheduler) Shutdown() {}
7279
7380 type JetstreamConsumer struct {
74- client *jsc.Client
75- logger *slog.Logger
76- sched *jetstreamScheduler
81+ client *jsc.Client
82+ logger *slog.Logger
83+ sched *jetstreamScheduler
84+ cursorStore CursorStore
85+ rewind time.Duration
7786 }
7887
79-func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slog.Logger) *JetstreamConsumer {
88+func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slog.Logger, cursorStore CursorStore) *JetstreamConsumer {
8089 sched := &jetstreamScheduler{
81- handler: handler,
82- logger: logger,
90+ handler: handler,
91+ logger: logger,
92+ cursorStore: cursorStore,
8393 }
8494
8595 wsURL := jetstreamURL
@@ -111,20 +121,32 @@ func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slo
111121 return nil
112122 }
113123
124+ rewind := 5 * time.Second
125+ if cursorStore == nil {
126+ rewind = 0
127+ }
128+
114129 return &JetstreamConsumer{
115- client: c,
116- logger: logger,
117- sched: sched,
130+ client: c,
131+ logger: logger,
132+ sched: sched,
133+ cursorStore: cursorStore,
134+ rewind: rewind,
118135 }
119136 }
120137
121138 func (jc *JetstreamConsumer) Start(ctx context.Context) error {
122139 for {
123- cursor := jc.sched.cursor.Load()
124140 var cursorPtr *int64
125- if cursor > 0 {
126- adjusted := cursor - int64(5*time.Second/time.Microsecond)
127- cursorPtr = &adjusted
141+ if jc.cursorStore != nil {
142+ cur, err := jc.cursorStore.LoadCursor(ctx)
143+ if err != nil {
144+ jc.logger.Warn("failed to load cursor, starting from now", "error", err)
145+ } else if cur != nil {
146+ rewound := max(*cur-int64(jc.rewind/time.Microsecond), 0)
147+ cursorPtr = &rewound
148+ jc.logger.Info("resuming jetstream", "cursor_us", *cur, "rewound_us", rewound)
149+ }
128150 }
129151
130152 err := jc.client.ConnectAndRead(ctx, cursorPtr)
@@ -10,11 +10,16 @@ import (
10 10
11 jsc "github.com/bluesky-social/jetstream/pkg/client"11 jsc "github.com/bluesky-social/jetstream/pkg/client"
12 "github.com/bluesky-social/jetstream/pkg/models"12 "github.com/bluesky-social/jetstream/pkg/models"
13- "go.uber.org/atomic"
14 13
15 "pkg.rbrt.fr/glean/internal/metrics"14 "pkg.rbrt.fr/glean/internal/metrics"
16 )15 )
17 16
17+// CursorStore stores the Jetstream cursor.
18+type CursorStore interface {
19+ LoadCursor(ctx context.Context) (*int64, error)
20+ SaveCursor(ctx context.Context, cursor int64) error
21+}
22+
18 type Event struct {23 type Event struct {
19 Type string24 Type string
20 DID string25 DID string
@@ -28,14 +33,16 @@ type Event struct {
28 type EventHandler func(ctx context.Context, event *Event) error33 type EventHandler func(ctx context.Context, event *Event) error
29 34
30 type jetstreamScheduler struct {35 type jetstreamScheduler struct {
31- handler EventHandler36+ handler EventHandler
32- logger *slog.Logger37+ logger *slog.Logger
33- cursor atomic.Int6438+ cursorStore CursorStore
34 }39 }
35 40
36 func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.Event) error {41 func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.Event) error {
37 if evt.TimeUS > 0 {42 if evt.TimeUS > 0 {
38- s.cursor.Store(evt.TimeUS)43+ if err := s.cursorStore.SaveCursor(ctx, evt.TimeUS); err != nil {
44+ s.logger.Warn("failed to save cursor", "error", err)
45+ }
39 }46 }
40 47
41 if evt.Kind != models.EventKindCommit || evt.Commit == nil {48 if evt.Kind != models.EventKindCommit || evt.Commit == nil {
@@ -71,15 +78,18 @@ func (s *jetstreamScheduler) AddWork(ctx context.Context, _ string, evt *models.
71 func (s *jetstreamScheduler) Shutdown() {}78 func (s *jetstreamScheduler) Shutdown() {}
72 79
73 type JetstreamConsumer struct {80 type JetstreamConsumer struct {
74- client *jsc.Client81+ client *jsc.Client
75- logger *slog.Logger82+ logger *slog.Logger
76- sched *jetstreamScheduler83+ sched *jetstreamScheduler
84+ cursorStore CursorStore
85+ rewind time.Duration
77 }86 }
78 87
79-func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slog.Logger) *JetstreamConsumer {88+func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slog.Logger, cursorStore CursorStore) *JetstreamConsumer {
80 sched := &jetstreamScheduler{89 sched := &jetstreamScheduler{
81- handler: handler,90+ handler: handler,
82- logger: logger,91+ logger: logger,
92+ cursorStore: cursorStore,
83 }93 }
84 94
85 wsURL := jetstreamURL95 wsURL := jetstreamURL
@@ -111,20 +121,32 @@ func NewJetstreamConsumer(jetstreamURL string, handler EventHandler, logger *slo
111 return nil121 return nil
112 }122 }
113 123
124+ rewind := 5 * time.Second
125+ if cursorStore == nil {
126+ rewind = 0
127+ }
128+
114 return &JetstreamConsumer{129 return &JetstreamConsumer{
115- client: c,130+ client: c,
116- logger: logger,131+ logger: logger,
117- sched: sched,132+ sched: sched,
133+ cursorStore: cursorStore,
134+ rewind: rewind,
118 }135 }
119 }136 }
120 137
121 func (jc *JetstreamConsumer) Start(ctx context.Context) error {138 func (jc *JetstreamConsumer) Start(ctx context.Context) error {
122 for {139 for {
123- cursor := jc.sched.cursor.Load()
124 var cursorPtr *int64140 var cursorPtr *int64
125- if cursor > 0 {141+ if jc.cursorStore != nil {
126- adjusted := cursor - int64(5*time.Second/time.Microsecond)142+ cur, err := jc.cursorStore.LoadCursor(ctx)
127- cursorPtr = &adjusted143+ if err != nil {
144+ jc.logger.Warn("failed to load cursor, starting from now", "error", err)
145+ } else if cur != nil {
146+ rewound := max(*cur-int64(jc.rewind/time.Microsecond), 0)
147+ cursorPtr = &rewound
148+ jc.logger.Info("resuming jetstream", "cursor_us", *cur, "rewound_us", rewound)
149+ }
128 }150 }
129 151
130 err := jc.client.ConnectAndRead(ctx, cursorPtr)152 err := jc.client.ConnectAndRead(ctx, cursorPtr)
added internal/db/cursor_store.go +34 -0
new file mode 100644
@@ -0,0 +1,34 @@
1+package db
2+
3+import (
4+ "context"
5+ "database/sql"
6+)
7+
8+type DBCursorStore struct {
9+ db *DB
10+}
11+
12+func NewCursorStore(db *DB) *DBCursorStore {
13+ return &DBCursorStore{db: db}
14+}
15+
16+func (s *DBCursorStore) LoadCursor(ctx context.Context) (*int64, error) {
17+ var cursor int64
18+ err := s.db.QueryRowContext(ctx, "SELECT cursor_us FROM jetstream_cursor WHERE id = 1").Scan(&cursor)
19+ if err == sql.ErrNoRows {
20+ return nil, nil
21+ }
22+ if err != nil {
23+ return nil, err
24+ }
25+ return &cursor, nil
26+}
27+
28+func (s *DBCursorStore) SaveCursor(ctx context.Context, cursor int64) error {
29+ _, err := s.db.ExecContext(ctx,
30+ "INSERT INTO jetstream_cursor (id, cursor_us) VALUES (1, ?) ON CONFLICT(id) DO UPDATE SET cursor_us = excluded.cursor_us",
31+ cursor,
32+ )
33+ return err
34+}
new file mode 100644
@@ -0,0 +1,34 @@
1+package db
2+
3+import (
4+ "context"
5+ "database/sql"
6+)
7+
8+type DBCursorStore struct {
9+ db *DB
10+}
11+
12+func NewCursorStore(db *DB) *DBCursorStore {
13+ return &DBCursorStore{db: db}
14+}
15+
16+func (s *DBCursorStore) LoadCursor(ctx context.Context) (*int64, error) {
17+ var cursor int64
18+ err := s.db.QueryRowContext(ctx, "SELECT cursor_us FROM jetstream_cursor WHERE id = 1").Scan(&cursor)
19+ if err == sql.ErrNoRows {
20+ return nil, nil
21+ }
22+ if err != nil {
23+ return nil, err
24+ }
25+ return &cursor, nil
26+}
27+
28+func (s *DBCursorStore) SaveCursor(ctx context.Context, cursor int64) error {
29+ _, err := s.db.ExecContext(ctx,
30+ "INSERT INTO jetstream_cursor (id, cursor_us) VALUES (1, ?) ON CONFLICT(id) DO UPDATE SET cursor_us = excluded.cursor_us",
31+ cursor,
32+ )
33+ return err
34+}
added internal/db/cursor_store_test.go +46 -0
new file mode 100644
@@ -0,0 +1,46 @@
1+package db
2+
3+import (
4+ "context"
5+ "testing"
6+
7+ "gotest.tools/v3/assert"
8+)
9+
10+func TestCursorStore_LoadEmpty(t *testing.T) {
11+ ctx := context.Background()
12+ dbs := setupTestDB(t)
13+ store := dbs.CursorStore()
14+
15+ cur, err := store.LoadCursor(ctx)
16+ assert.NilError(t, err)
17+ assert.Assert(t, cur == nil, "expected nil cursor for empty table")
18+}
19+
20+func TestCursorStore_SaveAndLoad(t *testing.T) {
21+ ctx := context.Background()
22+ dbs := setupTestDB(t)
23+ store := dbs.CursorStore()
24+
25+ err := store.SaveCursor(ctx, 1234567890)
26+ assert.NilError(t, err)
27+
28+ cur, err := store.LoadCursor(ctx)
29+ assert.NilError(t, err)
30+ assert.Assert(t, cur != nil)
31+ assert.Equal(t, *cur, int64(1234567890))
32+}
33+
34+func TestCursorStore_Update(t *testing.T) {
35+ ctx := context.Background()
36+ dbs := setupTestDB(t)
37+ store := dbs.CursorStore()
38+
39+ assert.NilError(t, store.SaveCursor(ctx, 100))
40+ assert.NilError(t, store.SaveCursor(ctx, 200))
41+
42+ cur, err := store.LoadCursor(ctx)
43+ assert.NilError(t, err)
44+ assert.Assert(t, cur != nil)
45+ assert.Equal(t, *cur, int64(200))
46+}
new file mode 100644
@@ -0,0 +1,46 @@
1+package db
2+
3+import (
4+ "context"
5+ "testing"
6+
7+ "gotest.tools/v3/assert"
8+)
9+
10+func TestCursorStore_LoadEmpty(t *testing.T) {
11+ ctx := context.Background()
12+ dbs := setupTestDB(t)
13+ store := dbs.CursorStore()
14+
15+ cur, err := store.LoadCursor(ctx)
16+ assert.NilError(t, err)
17+ assert.Assert(t, cur == nil, "expected nil cursor for empty table")
18+}
19+
20+func TestCursorStore_SaveAndLoad(t *testing.T) {
21+ ctx := context.Background()
22+ dbs := setupTestDB(t)
23+ store := dbs.CursorStore()
24+
25+ err := store.SaveCursor(ctx, 1234567890)
26+ assert.NilError(t, err)
27+
28+ cur, err := store.LoadCursor(ctx)
29+ assert.NilError(t, err)
30+ assert.Assert(t, cur != nil)
31+ assert.Equal(t, *cur, int64(1234567890))
32+}
33+
34+func TestCursorStore_Update(t *testing.T) {
35+ ctx := context.Background()
36+ dbs := setupTestDB(t)
37+ store := dbs.CursorStore()
38+
39+ assert.NilError(t, store.SaveCursor(ctx, 100))
40+ assert.NilError(t, store.SaveCursor(ctx, 200))
41+
42+ cur, err := store.LoadCursor(ctx)
43+ assert.NilError(t, err)
44+ assert.Assert(t, cur != nil)
45+ assert.Equal(t, *cur, int64(200))
46+}
modified internal/db/db.go +9 -0
@@ -175,6 +175,10 @@ func (s *Store) SQLDB() *sql.DB {
175175 return s.db.DB
176176 }
177177
178+func (s *Store) CursorStore() *DBCursorStore {
179+ return NewCursorStore(s.db)
180+}
181+
178182 func initUsersSchema(db *DB) error {
179183 for _, s := range usersSchema {
180184 if _, err := db.Exec(s); err != nil {
@@ -265,6 +269,11 @@ var usersSchema = []string{
265269 `CREATE INDEX IF NOT EXISTS idx_dismissed_user_type ON dismissed_recommendations(user_did, target_type)`,
266270 `CREATE INDEX IF NOT EXISTS idx_impressions_user_unacted ON recommendation_impressions(user_did, acted, shown_count)`,
267271 `CREATE INDEX IF NOT EXISTS idx_impressions_last_shown ON recommendation_impressions(last_shown_at)`,
272+
273+ `CREATE TABLE IF NOT EXISTS jetstream_cursor (
274+ id INTEGER PRIMARY KEY CHECK(id = 1),
275+ cursor_us INTEGER NOT NULL
276+ )`,
268277 }
269278
270279 var articlesSchema = []string{
@@ -175,6 +175,10 @@ func (s *Store) SQLDB() *sql.DB {
175 return s.db.DB175 return s.db.DB
176 }176 }
177 177
178+func (s *Store) CursorStore() *DBCursorStore {
179+ return NewCursorStore(s.db)
180+}
181+
178 func initUsersSchema(db *DB) error {182 func initUsersSchema(db *DB) error {
179 for _, s := range usersSchema {183 for _, s := range usersSchema {
180 if _, err := db.Exec(s); err != nil {184 if _, err := db.Exec(s); err != nil {
@@ -265,6 +269,11 @@ var usersSchema = []string{
265 `CREATE INDEX IF NOT EXISTS idx_dismissed_user_type ON dismissed_recommendations(user_did, target_type)`,269 `CREATE INDEX IF NOT EXISTS idx_dismissed_user_type ON dismissed_recommendations(user_did, target_type)`,
266 `CREATE INDEX IF NOT EXISTS idx_impressions_user_unacted ON recommendation_impressions(user_did, acted, shown_count)`,270 `CREATE INDEX IF NOT EXISTS idx_impressions_user_unacted ON recommendation_impressions(user_did, acted, shown_count)`,
267 `CREATE INDEX IF NOT EXISTS idx_impressions_last_shown ON recommendation_impressions(last_shown_at)`,271 `CREATE INDEX IF NOT EXISTS idx_impressions_last_shown ON recommendation_impressions(last_shown_at)`,
272+
273+ `CREATE TABLE IF NOT EXISTS jetstream_cursor (
274+ id INTEGER PRIMARY KEY CHECK(id = 1),
275+ cursor_us INTEGER NOT NULL
276+ )`,
268 }277 }
269 278
270 var articlesSchema = []string{279 var articlesSchema = []string{
modified internal/db/migrations.go +17 -1
@@ -13,7 +13,7 @@ func init() {
1313 }
1414
1515 // SchemaVersion must be incremented each time a migration is added to the migrations slice (used so that fresh dbs skip running migrations).
16-const SchemaVersion = 3
16+const SchemaVersion = 4
1717
1818 type migration struct {
1919 id int
@@ -37,6 +37,11 @@ var migrations = []migration{
3737 name: "article_language_user_languages",
3838 run: migrateArticleLanguageUserLanguages,
3939 },
40+ {
41+ id: 4,
42+ name: "jetstream_cursor",
43+ run: migrateJetstreamCursor,
44+ },
4045 }
4146
4247 func runMigrations(db *DB) error {
@@ -213,3 +218,14 @@ func migrateArticleLanguageUserLanguages(db *DB) error {
213218
214219 return nil
215220 }
221+
222+func migrateJetstreamCursor(db *DB) error {
223+ _, err := db.Exec(`CREATE TABLE IF NOT EXISTS jetstream_cursor (
224+ id INTEGER PRIMARY KEY CHECK(id = 1),
225+ cursor_us INTEGER NOT NULL
226+ )`)
227+ if err != nil {
228+ return fmt.Errorf("create jetstream_cursor: %w", err)
229+ }
230+ return nil
231+}
@@ -13,7 +13,7 @@ func init() {
13 }13 }
14 14
15 // SchemaVersion must be incremented each time a migration is added to the migrations slice (used so that fresh dbs skip running migrations).15 // SchemaVersion must be incremented each time a migration is added to the migrations slice (used so that fresh dbs skip running migrations).
16-const SchemaVersion = 316+const SchemaVersion = 4
17 17
18 type migration struct {18 type migration struct {
19 id int19 id int
@@ -37,6 +37,11 @@ var migrations = []migration{
37 name: "article_language_user_languages",37 name: "article_language_user_languages",
38 run: migrateArticleLanguageUserLanguages,38 run: migrateArticleLanguageUserLanguages,
39 },39 },
40+ {
41+ id: 4,
42+ name: "jetstream_cursor",
43+ run: migrateJetstreamCursor,
44+ },
40 }45 }
41 46
42 func runMigrations(db *DB) error {47 func runMigrations(db *DB) error {
@@ -213,3 +218,14 @@ func migrateArticleLanguageUserLanguages(db *DB) error {
213 218
214 return nil219 return nil
215 }220 }
221+
222+func migrateJetstreamCursor(db *DB) error {
223+ _, err := db.Exec(`CREATE TABLE IF NOT EXISTS jetstream_cursor (
224+ id INTEGER PRIMARY KEY CHECK(id = 1),
225+ cursor_us INTEGER NOT NULL
226+ )`)
227+ if err != nil {
228+ return fmt.Errorf("create jetstream_cursor: %w", err)
229+ }
230+ return nil
231+}
modified main.go +1 -1
@@ -91,7 +91,7 @@ func main() {
9191 cron := cluster.NewCron(engine, *clusterInterval, logger)
9292
9393 handler := atproto.NewStreamDBHandler(dbs.Articles, dbs.Users, logger)
94- jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger)
94+ jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger, dbs.CursorStore())
9595
9696 ctx, cancel := context.WithCancel(context.Background())
9797 defer cancel()
@@ -91,7 +91,7 @@ func main() {
91 cron := cluster.NewCron(engine, *clusterInterval, logger)91 cron := cluster.NewCron(engine, *clusterInterval, logger)
92 92
93 handler := atproto.NewStreamDBHandler(dbs.Articles, dbs.Users, logger)93 handler := atproto.NewStreamDBHandler(dbs.Articles, dbs.Users, logger)
94- jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger)94+ jetstream := atproto.NewJetstreamConsumer(*jetstreamURL, handler.Handle, logger, dbs.CursorStore())
95 95
96 ctx, cancel := context.WithCancel(context.Background())96 ctx, cancel := context.WithCancel(context.Background())
97 defer cancel()97 defer cancel()