diff --git a/bgs/admin.go b/bgs/admin.go new file mode 100644 index 00000000..0300d806 --- /dev/null +++ b/bgs/admin.go @@ -0,0 +1,11 @@ +package bgs + +import "github.com/labstack/echo/v4" + +func (bgs *BGS) handleAdminDeleteRecord(e echo.Context) error { + panic("TODO") +} + +func (bgs *BGS) handleAdminBlockRepoStream(e echo.Context) error { + panic("TODO") +} diff --git a/bgs/bgs.go b/bgs/bgs.go index 34bd8adc..5a27fac5 100644 --- a/bgs/bgs.go +++ b/bgs/bgs.go @@ -60,6 +60,7 @@ type BGS struct { func NewBGS(db *gorm.DB, ix *indexer.Indexer, repoman *repomgr.RepoManager, evtman *events.EventManager, didr plc.DidResolver, blobs blobs.BlobStore, ssl bool) (*BGS, error) { db.AutoMigrate(User{}) + db.AutoMigrate(AuthToken{}) db.AutoMigrate(models.PDS{}) bgs := &BGS{ @@ -181,7 +182,6 @@ func (bgs *BGS) Start(listen string) error { // TODO: this API is temporary until we formalize what we want here e.GET("/xrpc/com.atproto.sync.subscribeRepos", bgs.EventsHandler) - e.GET("/xrpc/com.atproto.sync.getCheckout", bgs.HandleComAtprotoSyncGetCheckout) e.GET("/xrpc/com.atproto.sync.getCommitPath", bgs.HandleComAtprotoSyncGetCommitPath) e.GET("/xrpc/com.atproto.sync.getHead", bgs.HandleComAtprotoSyncGetHead) @@ -192,6 +192,9 @@ func (bgs *BGS) Start(listen string) error { e.GET("/xrpc/com.atproto.sync.notifyOfUpdate", bgs.HandleComAtprotoSyncNotifyOfUpdate) e.GET("/xrpc/_health", bgs.HandleHealthCheck) + admin := e.Group("/admin", bgs.checkAdminAuth) + admin.POST("/deleteRecord", bgs.handleAdminDeleteRecord) + return e.Start(listen) } @@ -209,6 +212,62 @@ func (bgs *BGS) HandleHealthCheck(c echo.Context) error { } } +type AuthToken struct { + gorm.Model + Token string `gorm:"index"` +} + +func (bgs *BGS) lookupAdminToken(tok string) (bool, error) { + var at AuthToken + if err := bgs.db.Find(&at, "token = ?", tok).Error; err != nil { + return false, err + } + + if at.ID == 0 { + return false, nil + } + + return true, nil +} + +func (bgs *BGS) CreateAdminToken(tok string) error { + exists, err := bgs.lookupAdminToken(tok) + if err != nil { + return err + } + + if exists { + return nil + } + + return bgs.db.Create(&AuthToken{ + Token: tok, + }).Error +} + +func (bgs *BGS) checkAdminAuth(next echo.HandlerFunc) echo.HandlerFunc { + return func(e echo.Context) error { + authheader := e.Request().Header.Get("Authorization") + pref := "Bearer " + if !strings.HasPrefix(authheader, pref) { + return echo.ErrForbidden + } + + token := authheader[len(pref):] + + exists, err := bgs.lookupAdminToken(token) + if err != nil { + return err + } + + if !exists { + return echo.ErrForbidden + } + + return next(e) + } +} + type User struct { ID bsutil.Uid `gorm:"primarykey"` CreatedAt time.Time diff --git a/cmd/bigsky/main.go b/cmd/bigsky/main.go index 628c36e7..89378699 100644 --- a/cmd/bigsky/main.go +++ b/cmd/bigsky/main.go @@ -103,6 +103,10 @@ func run(args []string) { &cli.StringFlag{ Name: "disk-blob-store", }, + &cli.StringFlag{ + Name: "admin-key", + EnvVars: []string{"BGS_ADMIN_KEY"}, + }, } app.Action = func(cctx *cli.Context) error { @@ -201,6 +205,12 @@ func run(args []string) { return err } + if tok := cctx.String("admin-key"); tok != "" { + if err := bgs.CreateAdminToken(tok); err != nil { + return fmt.Errorf("failed to set up admin token: %w", err) + } + } + // set up pprof endpoint go func() { if err := bgs.StartDebug(cctx.String("debug-listen")); err != nil { -- 2.51.2 From 5fba912435cb8fa70c8617857a3d3d540b63562b Mon Sep 17 00:00:00 2001 From: whyrusleeping Date: Thu, 27 Apr 2023 11:17:51 -0700 Subject: [PATCH 2/4] add admin route to disable new subscriptions --- bgs/admin.go | 27 ++++++++++++++++++++++- bgs/bgs.go | 7 +++++- bgs/fedmgr.go | 53 +++++++++++++++++++++++++++++++++++++++++++--- indexer/indexer.go | 33 ++++------------------------- 4 files changed, 86 insertions(+), 34 deletions(-) diff --git a/bgs/admin.go b/bgs/admin.go index 0300d806..cbce9fc6 100644 --- a/bgs/admin.go +++ b/bgs/admin.go @@ -1,11 +1,36 @@ package bgs -import "github.com/labstack/echo/v4" +import ( + "strconv" + + "github.com/bluesky-social/indigo/util" + "github.com/labstack/echo/v4" +) func (bgs *BGS) handleAdminDeleteRecord(e echo.Context) error { + puri, err := util.ParseAtUri(e.QueryParam("uri")) + if err != nil { + return err + } + + _ = puri + panic("TODO") } func (bgs *BGS) handleAdminBlockRepoStream(e echo.Context) error { panic("TODO") } + +func (bgs *BGS) handleAdminDisableNewSlurps(e echo.Context) error { + enabled, err := strconv.ParseBool(e.QueryParam("enabled")) + if err != nil { + return err + } + + return bgs.slurper.SetNewSubsDisabled(!enabled) +} + +func (bgs *BGS) handleAdminTakedownRepo(e echo.Context) error { + panic("TODO") +} diff --git a/bgs/bgs.go b/bgs/bgs.go index 5a27fac5..b7c08356 100644 --- a/bgs/bgs.go +++ b/bgs/bgs.go @@ -74,7 +74,12 @@ func NewBGS(db *gorm.DB, ix *indexer.Indexer, repoman *repomgr.RepoManager, evtm } ix.CreateExternalUser = bgs.createExternalUser - bgs.slurper = NewSlurper(db, bgs.handleFedEvent, ssl) + s, err := NewSlurper(db, bgs.handleFedEvent, ssl) + if err != nil { + return nil, err + } + + bgs.slurper = s if err := bgs.slurper.RestartAll(); err != nil { return nil, err diff --git a/bgs/fedmgr.go b/bgs/fedmgr.go index e884ffde..a2bbdbee 100644 --- a/bgs/fedmgr.go +++ b/bgs/fedmgr.go @@ -27,22 +27,69 @@ type Slurper struct { lk sync.Mutex active map[string]*models.PDS + newSubsDisabled bool + ssl bool } -func NewSlurper(db *gorm.DB, cb IndexCallback, ssl bool) *Slurper { - return &Slurper{ +func NewSlurper(db *gorm.DB, cb IndexCallback, ssl bool) (*Slurper, error) { + s := &Slurper{ cb: cb, db: db, active: make(map[string]*models.PDS), ssl: ssl, } + if err := s.loadConfig(); err != nil { + return nil, err + } + + return s, nil +} + +func (s *Slurper) loadConfig() error { + var sc SlurpConfig + if err := s.db.Find(&sc).Error; err != nil { + return err + } + + if sc.ID == 0 { + if err := s.db.Create(&SlurpConfig{}).Error; err != nil { + return err + } + } + + s.newSubsDisabled = sc.NewSubsDisabled + + return nil +} + +type SlurpConfig struct { + gorm.Model + + NewSubsDisabled bool +} + +func (s *Slurper) SetNewSubsDisabled(dis bool) error { + s.lk.Lock() + defer s.lk.Unlock() + + if err := s.db.Model(SlurpConfig{}).Where("id = 1").Update("new_subs_disabled", dis).Error; err != nil { + return err + } + + s.newSubsDisabled = dis + return nil } +var ErrNewSubsDisabled = fmt.Errorf("new subscriptions temporarily disabled") + func (s *Slurper) SubscribeToPds(ctx context.Context, host string, reg bool) error { // TODO: for performance, lock on the hostname instead of global s.lk.Lock() defer s.lk.Unlock() + if s.newSubsDisabled { + return ErrNewSubsDisabled + } _, ok := s.active[host] if ok { @@ -87,7 +134,7 @@ func (s *Slurper) RestartAll() error { defer s.lk.Unlock() var all []models.PDS - if err := s.db.Find(&all).Error; err != nil { + if err := s.db.Find(&all, "registered = true").Error; err != nil { return err } diff --git a/indexer/indexer.go b/indexer/indexer.go index 4257d0de..067ec627 100644 --- a/indexer/indexer.go +++ b/indexer/indexer.go @@ -5,7 +5,6 @@ import ( "context" "errors" "fmt" - "strings" "time" comatproto "github.com/bluesky-social/indigo/api/atproto" @@ -310,7 +309,7 @@ func (ix *Indexer) handleRecordCreate(ctx context.Context, evt *repomgr.RepoEven } func (ix *Indexer) crawlAtUriRef(ctx context.Context, uri string) error { - puri, err := parseAtUri(uri) + puri, err := util.ParseAtUri(uri) if err != nil { return err } else { @@ -538,7 +537,7 @@ func (ix *Indexer) handleRecordUpdate(ctx context.Context, evt *repomgr.RepoEven } func (ix *Indexer) GetPostOrMissing(ctx context.Context, uri string) (*models.FeedPost, error) { - puri, err := parseAtUri(uri) + puri, err := util.ParseAtUri(uri) if err != nil { return nil, err } @@ -643,7 +642,7 @@ func (ix *Indexer) GetUserOrMissing(ctx context.Context, did string) (*models.Ac return ix.createMissingUserRecord(ctx, did) } -func (ix *Indexer) createMissingPostRecord(ctx context.Context, puri *parsedUri) (*models.FeedPost, error) { +func (ix *Indexer) createMissingPostRecord(ctx context.Context, puri *util.ParsedUri) (*models.FeedPost, error) { log.Warn("creating missing post record") ai, err := ix.GetUserOrMissing(ctx, puri.Did) if err != nil { @@ -793,7 +792,7 @@ func isNotFound(err error) bool { } func (ix *Indexer) GetPost(ctx context.Context, uri string) (*models.FeedPost, error) { - puri, err := parseAtUri(uri) + puri, err := util.ParseAtUri(uri) if err != nil { return nil, err } @@ -806,30 +805,6 @@ func (ix *Indexer) GetPost(ctx context.Context, uri string) (*models.FeedPost, e return &post, nil } -type parsedUri struct { - Did string - Collection string - Rkey string -} - -func parseAtUri(uri string) (*parsedUri, error) { - if !strings.HasPrefix(uri, "at://") { - return nil, fmt.Errorf("AT uris must be prefixed with 'at://'") - } - - trimmed := strings.TrimPrefix(uri, "at://") - parts := strings.Split(trimmed, "/") - if len(parts) != 3 { - return nil, fmt.Errorf("AT uris must have three parts: did, collection, tid") - } - - return &parsedUri{ - Did: parts[0], - Collection: parts[1], - Rkey: parts[2], - }, nil -} - // TODO: since this function is the only place we depend on the repomanager, i wonder if this should be wired some other way? func (ix *Indexer) FetchAndIndexRepo(ctx context.Context, job *crawlWork) error { ctx, span := otel.Tracer("indexer").Start(ctx, "FetchAndIndexRepo") -- 2.51.2 From 119833a25ecc408339e7416ff9528fc4c5429e7f Mon Sep 17 00:00:00 2001 From: whyrusleeping Date: Tue, 2 May 2023 11:25:37 -0700 Subject: [PATCH 3/4] fix build --- bgs/fedmgr.go | 1 + labeler/service.go | 5 ++++- 2 files changed, 5 insertions(+), 1 deletion(-) diff --git a/bgs/fedmgr.go b/bgs/fedmgr.go index a2bbdbee..cb4da324 100644 --- a/bgs/fedmgr.go +++ b/bgs/fedmgr.go @@ -33,6 +33,7 @@ type Slurper struct { } func NewSlurper(db *gorm.DB, cb IndexCallback, ssl bool) (*Slurper, error) { + db.AutoMigrate(&SlurpConfig{}) s := &Slurper{ cb: cb, db: db, diff --git a/labeler/service.go b/labeler/service.go index 91e69503..df2fd4bf 100644 --- a/labeler/service.go +++ b/labeler/service.go @@ -111,7 +111,10 @@ func NewServer(db *gorm.DB, cs *carstore.CarStore, repoUser RepoConfig, plcURL, log.Infof("found labelmaker repo: %s", head) } - slurp := bgs.NewSlurper(db, s.handleBgsRepoEvent, useWss) + slurp, err := bgs.NewSlurper(db, s.handleBgsRepoEvent, useWss) + if err != nil { + return nil, err + } s.bgsSlurper = slurp go evtmgr.Run() -- 2.51.2 From f037747885d78b93e910d1ba55a01f1b7e288c27 Mon Sep 17 00:00:00 2001 From: whyrusleeping Date: Tue, 2 May 2023 12:43:21 -0700 Subject: [PATCH 4/4] add missing file --- util/uri.go | 30 ++++++++++++++++++++++++++++++ 1 file changed, 30 insertions(+) create mode 100644 util/uri.go diff --git a/util/uri.go b/util/uri.go new file mode 100644 index 00000000..cfc6bacb --- /dev/null +++ b/util/uri.go @@ -0,0 +1,30 @@ +package util + +import ( + "fmt" + "strings" +) + +type ParsedUri struct { + Did string + Collection string + Rkey string +} + +func ParseAtUri(uri string) (*ParsedUri, error) { + if !strings.HasPrefix(uri, "at://") { + return nil, fmt.Errorf("AT uris must be prefixed with 'at://'") + } + + trimmed := strings.TrimPrefix(uri, "at://") + parts := strings.Split(trimmed, "/") + if len(parts) != 3 { + return nil, fmt.Errorf("AT uris must have three parts: did, collection, tid") + } + + return &ParsedUri{ + Did: parts[0], + Collection: parts[1], + Rkey: parts[2], + }, nil +}