From cff1285e349e046e499639f8d6f6ddc6894ff026 Mon Sep 17 00:00:00 2001 From: Will Date: Sat, 6 Jun 2026 15:37:44 +0100 Subject: [PATCH] event persisting using retryable sequencing increments on conflicts Signed-off-by: Will --- go.mod | 1 + go.sum | 2 + server/persist.go | 111 ++++++++++++++++++++++++---------------------- server/server.go | 9 ++-- 4 files changed, 68 insertions(+), 55 deletions(-) diff --git a/go.mod b/go.mod index 0b9e30f..f219ceb 100644 --- a/go.mod +++ b/go.mod @@ -4,6 +4,7 @@ go 1.25 require ( github.com/Azure/go-autorest/autorest/to v0.4.1 + github.com/avast/retry-go/v4 v4.7.0 github.com/aws/aws-sdk-go v1.55.7 github.com/bluesky-social/go-util v0.0.0-20251012040650-2ebbf57f5934 github.com/bluesky-social/indigo v0.0.0-20260203235305-a86f3ae1f8ec diff --git a/go.sum b/go.sum index 92e0743..18b8929 100644 --- a/go.sum +++ b/go.sum @@ -9,6 +9,8 @@ github.com/alexbrainman/goissue34681 v0.0.0-20191006012335-3fc7a47baff5 h1:iW0a5 github.com/alexbrainman/goissue34681 v0.0.0-20191006012335-3fc7a47baff5/go.mod h1:Y2QMoi1vgtOIfc+6DhrMOGkLoGzqSV2rKp4Sm+opsyA= github.com/antlr4-go/antlr/v4 v4.13.0 h1:lxCg3LAv+EUK6t1i0y1V6/SLeUi0eKEKdhQAlS8TVTI= github.com/antlr4-go/antlr/v4 v4.13.0/go.mod h1:pfChB/xh/Unjila75QW7+VU4TSnWnnk9UTnmpPaOR2g= +github.com/avast/retry-go/v4 v4.7.0 h1:yjDs35SlGvKwRNSykujfjdMxMhMQQM0TnIjJaHB+Zio= +github.com/avast/retry-go/v4 v4.7.0/go.mod h1:ZMPDa3sY2bKgpLtap9JRUgk2yTAba7cgiFhqxY2Sg6Q= github.com/aws/aws-sdk-go v1.55.7 h1:UJrkFq7es5CShfBwlWAC8DA077vp8PyVbQd3lqLiztE= github.com/aws/aws-sdk-go v1.55.7/go.mod h1:eRwEWoyTWFMVYVQzKMNHWP5/RV4xIUGMQfXQHfHkpNU= github.com/benbjohnson/clock v1.1.0/go.mod h1:J11/hYXuz8f4ySSvYwY0FKfm+ezbsZBKZxNJlLklBHA= diff --git a/server/persist.go b/server/persist.go index 0ff510d..700061f 100644 --- a/server/persist.go +++ b/server/persist.go @@ -3,10 +3,13 @@ package server import ( "bytes" "context" + "errors" "fmt" + "log/slog" "sync" "time" + "github.com/avast/retry-go/v4" "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/events" indigomodels "github.com/bluesky-social/indigo/models" @@ -19,8 +22,7 @@ import ( type DbPersister struct { Db *gorm.DB - Lk sync.Mutex - Seq int64 + Lk sync.Mutex Broadcast func(*events.XRPCStreamEvent) @@ -42,19 +44,6 @@ func NewDbPersister(db *gorm.DB, retention time.Duration) (*DbPersister, error) Retention: retention, } - // kind of hacky. we will try and get the latest one from the db, but if it doesn't exist...well we have a problem - // because the relay will already have _some_ value > 0 set as a cursor, we'll want to just set this to some high value - // we'll just grab a current unix timestamp and set that as the cursor - var lastEvent models.EventRecord - if err := db.Order("seq desc").Limit(1).First(&lastEvent).Error; err != nil { - if err != gorm.ErrRecordNotFound { - return nil, fmt.Errorf("failed to get last event seq: %w", err) - } - p.Seq = time.Now().Unix() - } else { - p.Seq = lastEvent.Seq - } - go p.cleanupRoutine() return p, nil @@ -68,48 +57,66 @@ func (p *DbPersister) Persist(ctx context.Context, e *events.XRPCStreamEvent) er p.Lk.Lock() defer p.Lk.Unlock() - p.Seq++ - seq := p.Seq - - var did string - var evtType string - - switch { - case e.RepoCommit != nil: - e.RepoCommit.Seq = seq - did = e.RepoCommit.Repo - evtType = "commit" - case e.RepoSync != nil: - e.RepoSync.Seq = seq - did = e.RepoSync.Did - evtType = "sync" - case e.RepoIdentity != nil: - e.RepoIdentity.Seq = seq - did = e.RepoIdentity.Did - evtType = "identity" - case e.RepoAccount != nil: - e.RepoAccount.Seq = seq - did = e.RepoAccount.Did - evtType = "account" - default: - return fmt.Errorf("unknown event type") + rec := &models.EventRecord{} + if err := p.Db.Order("seq desc").Limit(1).First(rec).Error; err != nil { + slog.Error("fetching most recent event record", "error", err) + rec.Seq = time.Now().Unix() } - data, err := serializeEvent(e) - if err != nil { - return fmt.Errorf("failed to serialize event: %w", err) + // if the error on inserting the event record is a constraint error, it means that + // another event has been stored since getting the last sequence number. In that case + // retry which will increase the sequence number again and hopefully work. Any others + // can error out. + retryIfFunc := func(err error) bool { + return errors.Is(err, gorm.ErrDuplicatedKey) } + err := retry.Do(func() error { + rec.Seq++ + + var did string + var evtType string + + switch { + case e.RepoCommit != nil: + e.RepoCommit.Seq = rec.Seq + did = e.RepoCommit.Repo + evtType = "commit" + case e.RepoSync != nil: + e.RepoSync.Seq = rec.Seq + did = e.RepoSync.Did + evtType = "sync" + case e.RepoIdentity != nil: + e.RepoIdentity.Seq = rec.Seq + did = e.RepoIdentity.Did + evtType = "identity" + case e.RepoAccount != nil: + e.RepoAccount.Seq = rec.Seq + did = e.RepoAccount.Did + evtType = "account" + default: + return fmt.Errorf("unknown event type") + } - rec := &models.EventRecord{ - Seq: seq, - CreatedAt: time.Now(), - Did: did, - Type: evtType, - Data: data, - } + data, err := serializeEvent(e) + if err != nil { + return fmt.Errorf("failed to serialize event: %w", err) + } + + rec.CreatedAt = time.Now() + rec.Did = did + rec.Type = evtType + rec.Data = data - if err := p.Db.Create(rec).Error; err != nil { - return fmt.Errorf("failed to persist event: %w", err) + err = p.Db.Create(rec).Error + if err != nil { + return fmt.Errorf("failed to persist event: %w", err) + } + + return nil + }, retry.RetryIf(retryIfFunc)) + + if err != nil { + return err } if p.Broadcast != nil { diff --git a/server/server.go b/server/server.go index 79cebb9..898de9a 100644 --- a/server/server.go +++ b/server/server.go @@ -348,13 +348,16 @@ func New(args *Args) (*Server, error) { } var gdb *gorm.DB + gormCfg := gorm.Config{ + TranslateError: true, + } var err error switch dbType { case "postgres": if args.DatabaseURL == "" { return nil, fmt.Errorf("database-url must be set when using postgres") } - gdb, err = gorm.Open(postgres.Open(args.DatabaseURL), &gorm.Config{}) + gdb, err = gorm.Open(postgres.Open(args.DatabaseURL), &gormCfg) if err != nil { return nil, fmt.Errorf("failed to connect to postgres: %w", err) } @@ -366,14 +369,14 @@ func New(args *Args) (*Server, error) { db, err := sql.Open("libsql", fmt.Sprintf("%s?authToken=%s", primaryUrl, authToken)) gdb, err = gorm.Open(sqlite.New(sqlite.Config{ Conn: db, - }), &gorm.Config{}) + }), &gormCfg) if err != nil { return nil, fmt.Errorf("failed to connect to postgres: %w", err) } logger.Info("connected to PostgreSQL database") default: - gdb, err = gorm.Open(sqlite.Open(args.DbName), &gorm.Config{}) + gdb, err = gorm.Open(sqlite.Open(args.DbName), &gormCfg) if err != nil { return nil, fmt.Errorf("failed to open sqlite database: %w", err) } -- 2.51.2