diff --git a/appview/firehose.go b/appview/firehose.go index 3603b10..f0dcde6 100644 --- a/appview/firehose.go +++ b/appview/firehose.go @@ -18,10 +18,15 @@ import ( "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" "github.com/gorilla/websocket" + "gorm.io/gorm/clause" "tangled.org/sparrowtek.com/effem-AppView/appview/database" "tangled.org/sparrowtek.com/effem-AppView/appview/metrics" ) +// firehoseCursorRowID is the singleton row holding the stream seq. We pin it +// to 1 so the persist path can run as a single upsert keyed on the PK. +const firehoseCursorRowID = 1 + const ( effemNSPrefix = "xyz.effem." @@ -375,13 +380,13 @@ func (srv *Server) loadFirehoseCursor(ctx context.Context) (int64, error) { } func (srv *Server) saveFirehoseCursor(ctx context.Context, seq int64) error { - db := srv.db.WithContext(ctx) - var cur database.FirehoseCursor - err := db.First(&cur).Error - if err != nil { - cur = database.FirehoseCursor{Seq: seq} - return db.Create(&cur).Error + cur := database.FirehoseCursor{ + ID: firehoseCursorRowID, + Seq: seq, + UpdatedAt: time.Now(), } - cur.Seq = seq - return db.Save(&cur).Error + return srv.db.WithContext(ctx).Clauses(clause.OnConflict{ + Columns: []clause.Column{{Name: "id"}}, + DoUpdates: clause.AssignmentColumns([]string{"seq", "updated_at"}), + }).Create(&cur).Error }