diff --git a/bgs/fedmgr.go b/bgs/fedmgr.go index c7dc0565..b079f708 100644 --- a/bgs/fedmgr.go +++ b/bgs/fedmgr.go @@ -10,7 +10,7 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/events" - "github.com/bluesky-social/indigo/events/autoscaling" + "github.com/bluesky-social/indigo/events/schedulers/autoscaling" "github.com/bluesky-social/indigo/models" "go.opentelemetry.io/otel" @@ -389,7 +389,7 @@ func (s *Slurper) handleConnection(ctx context.Context, host *models.PDS, con *w }, } - pool := autoscaling.NewConsumerPool(1, 360, con.RemoteAddr().String(), rsc.EventHandler) + pool := autoscaling.NewScheduler(1, 360, time.Second, con.RemoteAddr().String(), rsc.EventHandler) return events.HandleRepoStream(ctx, con, pool) } diff --git a/cmd/gosky/debug.go b/cmd/gosky/debug.go index ba1f7de8..007c5758 100644 --- a/cmd/gosky/debug.go +++ b/cmd/gosky/debug.go @@ -18,6 +18,7 @@ import ( "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/did" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/sequential" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" @@ -97,7 +98,8 @@ var inspectEventCmd = &cli.Command{ }, } - err = events.HandleRepoStream(ctx, con, &events.SequentialScheduler{rsc.EventHandler}) + seqScheduler := sequential.NewScheduler("debug-inspect-event", rsc.EventHandler) + err = events.HandleRepoStream(ctx, con, seqScheduler) if err != errFoundIt { return err } @@ -251,7 +253,8 @@ var debugStreamCmd = &cli.Command{ return fmt.Errorf("%s: %s", evt.Error, evt.Message) }, } - err = events.HandleRepoStream(ctx, con, &events.SequentialScheduler{rsc.EventHandler}) + seqScheduler := sequential.NewScheduler("debug-stream", rsc.EventHandler) + err = events.HandleRepoStream(ctx, con, seqScheduler) if err != nil { return err } @@ -371,7 +374,8 @@ var compareStreamsCmd = &cli.Command{ return fmt.Errorf("%s: %s", evt.Error, evt.Message) }, } - if err := events.HandleRepoStream(ctx, con, &events.SequentialScheduler{rsc.EventHandler}); err != nil { + seqScheduler := sequential.NewScheduler(fmt.Sprintf("debug-stream-%d", i+1), rsc.EventHandler) + if err := events.HandleRepoStream(ctx, con, seqScheduler); err != nil { log.Fatalf("HandleRepoStream failure on url%d: %s", i+1, err) } }(i, url) diff --git a/cmd/gosky/main.go b/cmd/gosky/main.go index 14687ebb..e435a0fb 100644 --- a/cmd/gosky/main.go +++ b/cmd/gosky/main.go @@ -20,6 +20,7 @@ import ( "github.com/bluesky-social/indigo/api/bsky" appbsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/sequential" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/util" @@ -1098,7 +1099,8 @@ var readRepoStreamCmd = &cli.Command{ return fmt.Errorf("error frame: %s: %s", errf.Error, errf.Message) }, } - return events.HandleRepoStream(ctx, con, &events.SequentialScheduler{rsc.EventHandler}) + seqScheduler := sequential.NewScheduler(con.RemoteAddr().String(), rsc.EventHandler) + return events.HandleRepoStream(ctx, con, seqScheduler) }, } diff --git a/cmd/gosky/streamdiff.go b/cmd/gosky/streamdiff.go index 3a245fd4..6ff50c9f 100644 --- a/cmd/gosky/streamdiff.go +++ b/cmd/gosky/streamdiff.go @@ -7,6 +7,7 @@ import ( comatproto "github.com/bluesky-social/indigo/api/atproto" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/sequential" "github.com/gorilla/websocket" cli "github.com/urfave/cli/v2" ) @@ -55,7 +56,8 @@ var streamCompareCmd = &cli.Command{ return fmt.Errorf("%s: %s", evt.Error, evt.Message) }, } - err = events.HandleRepoStream(ctx, cona, &events.SequentialScheduler{rsc.EventHandler}) + seqScheduler := sequential.NewScheduler("streamA", rsc.EventHandler) + err = events.HandleRepoStream(ctx, cona, seqScheduler) if err != nil { log.Errorf("stream A failed: %s", err) } @@ -77,9 +79,11 @@ var streamCompareCmd = &cli.Command{ return fmt.Errorf("%s: %s", evt.Error, evt.Message) }, } - err = events.HandleRepoStream(ctx, conb, &events.SequentialScheduler{rsc.EventHandler}) + + seqScheduler := sequential.NewScheduler("streamB", rsc.EventHandler) + err = events.HandleRepoStream(ctx, conb, seqScheduler) if err != nil { - log.Errorf("stream A failed: %s", err) + log.Errorf("stream B failed: %s", err) } }() diff --git a/cmd/sonar/main.go b/cmd/sonar/main.go index e8e1e350..05ce885f 100644 --- a/cmd/sonar/main.go +++ b/cmd/sonar/main.go @@ -13,6 +13,7 @@ import ( "time" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/autoscaling" "github.com/bluesky-social/indigo/sonar" "github.com/bluesky-social/indigo/util/version" "github.com/gorilla/websocket" @@ -108,7 +109,7 @@ func Sonar(cctx *cli.Context) error { wg := sync.WaitGroup{} - pool := events.NewConsumerPool(cctx.Int("worker-count"), cctx.Int("max-queue-size"), u.Host, s.HandleStreamEvent) + pool := autoscaling.NewScheduler(cctx.Int("worker-count"), cctx.Int("max-queue-size"), time.Second, u.Host, s.HandleStreamEvent) // Start a goroutine to manage the cursor file, saving the current cursor every 5 seconds. go func() { diff --git a/events/autoscaling/metrics.go b/events/autoscaling/metrics.go deleted file mode 100644 index 58b3c901..00000000 --- a/events/autoscaling/metrics.go +++ /dev/null @@ -1,26 +0,0 @@ -package autoscaling - -import ( - "github.com/prometheus/client_golang/prometheus" - "github.com/prometheus/client_golang/prometheus/promauto" -) - -var workItemsAdded = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "indigo_pool_work_items_added_total", - Help: "Total number of work items added to the consumer pool", -}, []string{"pool", "pool_type"}) - -var workItemsProcessed = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "indigo_pool_work_items_processed_total", - Help: "Total number of work items processed by the consumer pool", -}, []string{"pool", "pool_type"}) - -var workItemsActive = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "indigo_pool_work_items_active_total", - Help: "Total number of work items passed into a worker", -}, []string{"pool", "pool_type"}) - -var workersActive = promauto.NewGaugeVec(prometheus.GaugeOpts{ - Name: "indigo_pool_workers_active", - Help: "Number of workers currently active", -}, []string{"pool", "pool_type"}) diff --git a/events/events.go b/events/events.go index f7e5a72d..f727e290 100644 --- a/events/events.go +++ b/events/events.go @@ -17,6 +17,10 @@ import ( var log = logging.Logger("events") +type Scheduler interface { + AddWork(ctx context.Context, repo string, val *XRPCStreamEvent) error +} + type EventManager struct { subs []*Subscriber subsLk sync.Mutex diff --git a/events/metrics.go b/events/metrics.go index 972138f6..0788f2f7 100644 --- a/events/metrics.go +++ b/events/metrics.go @@ -15,21 +15,6 @@ var bytesFromStreamCounter = promauto.NewCounterVec(prometheus.CounterOpts{ Help: "Total bytes received from the stream", }, []string{"remote_addr"}) -var workItemsAdded = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "indigo_work_items_added_total", - Help: "Total number of work items added to the consumer pool", -}, []string{"pool"}) - -var workItemsProcessed = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "indigo_work_items_processed_total", - Help: "Total number of work items processed by the consumer pool", -}, []string{"pool"}) - -var workItemsActive = promauto.NewCounterVec(prometheus.CounterOpts{ - Name: "indigo_work_items_active_total", - Help: "Total number of work items passed into a worker", -}, []string{"pool"}) - var eventsEnqueued = promauto.NewCounterVec(prometheus.CounterOpts{ Name: "indigo_events_enqueued_for_broadcast_total", Help: "Total number of events enqueued to broadcast to subscribers", diff --git a/events/parallel.go b/events/parallel.go deleted file mode 100644 index 2471dd61..00000000 --- a/events/parallel.go +++ /dev/null @@ -1,110 +0,0 @@ -package events - -import ( - "context" - "sync" -) - -type Scheduler interface { - AddWork(ctx context.Context, repo string, val *XRPCStreamEvent) error -} - -type SequentialScheduler struct { - Do func(context.Context, *XRPCStreamEvent) error -} - -func (s *SequentialScheduler) AddWork(ctx context.Context, repo string, val *XRPCStreamEvent) error { - return s.Do(ctx, val) -} - -type ParallelConsumerPool struct { - maxConcurrency int - maxQueue int - - do func(context.Context, *XRPCStreamEvent) error - - feeder chan *consumerTask - - lk sync.Mutex - active map[string][]*consumerTask - - ident string -} - -func NewConsumerPool(maxC, maxQ int, ident string, do func(context.Context, *XRPCStreamEvent) error) *ParallelConsumerPool { - p := &ParallelConsumerPool{ - maxConcurrency: maxC, - maxQueue: maxQ, - - do: do, - - feeder: make(chan *consumerTask), - active: make(map[string][]*consumerTask), - - ident: ident, - } - - for i := 0; i < maxC; i++ { - go p.worker() - } - - return p -} - -type consumerTask struct { - repo string - val *XRPCStreamEvent -} - -func (p *ParallelConsumerPool) AddWork(ctx context.Context, repo string, val *XRPCStreamEvent) error { - workItemsAdded.WithLabelValues(p.ident).Inc() - t := &consumerTask{ - repo: repo, - val: val, - } - p.lk.Lock() - - a, ok := p.active[repo] - if ok { - p.active[repo] = append(a, t) - p.lk.Unlock() - return nil - } - - p.active[repo] = []*consumerTask{} - p.lk.Unlock() - - select { - case p.feeder <- t: - return nil - case <-ctx.Done(): - return ctx.Err() - } -} - -func (p *ParallelConsumerPool) worker() { - for work := range p.feeder { - for work != nil { - workItemsActive.WithLabelValues(p.ident).Inc() - if err := p.do(context.TODO(), work.val); err != nil { - log.Errorf("event handler failed: %s", err) - } - workItemsProcessed.WithLabelValues(p.ident).Inc() - - p.lk.Lock() - rem, ok := p.active[work.repo] - if !ok { - log.Errorf("should always have an 'active' entry if a worker is processing a job") - } - - if len(rem) == 0 { - delete(p.active, work.repo) - work = nil - } else { - work = rem[0] - p.active[work.repo] = rem[1:] - } - p.lk.Unlock() - } - } -} diff --git a/events/autoscaling/autoscaling.go b/events/schedulers/autoscaling/autoscaling.go similarity index 68% rename from events/autoscaling/autoscaling.go rename to events/schedulers/autoscaling/autoscaling.go index 3632d008..f3d5b3aa 100644 --- a/events/autoscaling/autoscaling.go +++ b/events/schedulers/autoscaling/autoscaling.go @@ -6,11 +6,13 @@ import ( "time" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers" "github.com/labstack/gommon/log" "github.com/prometheus/client_golang/prometheus" ) -type ConsumerPool struct { +// Scheduler is a scheduler that will scale up and down the number of workers based on the throughput of the workers. +type Scheduler struct { concurrency int maxConcurrency int @@ -27,14 +29,15 @@ type ConsumerPool struct { itemsAdded prometheus.Counter itemsProcessed prometheus.Counter itemsActive prometheus.Counter - workersAcrive prometheus.Gauge + workersActive prometheus.Gauge // autoscaling - throughputManager *ThroughputManager + throughputManager *ThroughputManager + autoscaleFrequency time.Duration } -func NewConsumerPool(concurrency, maxC int, ident string, do func(context.Context, *events.XRPCStreamEvent) error) *ConsumerPool { - p := &ConsumerPool{ +func NewScheduler(concurrency, maxC int, autoscaleFrequency time.Duration, ident string, do func(context.Context, *events.XRPCStreamEvent) error) *Scheduler { + p := &Scheduler{ concurrency: concurrency, maxConcurrency: maxC, @@ -45,14 +48,15 @@ func NewConsumerPool(concurrency, maxC int, ident string, do func(context.Contex ident: ident, - itemsAdded: workItemsAdded.WithLabelValues(ident, "autoscaling"), - itemsProcessed: workItemsProcessed.WithLabelValues(ident, "autoscaling"), - itemsActive: workItemsActive.WithLabelValues(ident, "autoscaling"), - workersAcrive: workersActive.WithLabelValues(ident, "autoscaling"), + itemsAdded: schedulers.WorkItemsAdded.WithLabelValues(ident, "autoscaling"), + itemsProcessed: schedulers.WorkItemsProcessed.WithLabelValues(ident, "autoscaling"), + itemsActive: schedulers.WorkItemsActive.WithLabelValues(ident, "autoscaling"), + workersActive: schedulers.WorkersActive.WithLabelValues(ident, "autoscaling"), // autoscaling // By default, the ThroughputManager will calculate the average throughput over the last 60 seconds. - throughputManager: NewThroughputManager(60), + throughputManager: NewThroughputManager(60), + autoscaleFrequency: autoscaleFrequency, } for i := 0; i < concurrency; i++ { @@ -65,9 +69,9 @@ func NewConsumerPool(concurrency, maxC int, ident string, do func(context.Contex } // Add autoscaling function -func (p *ConsumerPool) autoscale() { +func (p *Scheduler) autoscale() { p.throughputManager.Start() - tick := time.NewTicker(time.Second * 5) // adjust as needed + tick := time.NewTicker(p.autoscaleFrequency) for range tick.C { avg := p.throughputManager.AvgThroughput() if avg > float64(p.concurrency) && p.concurrency < p.maxConcurrency { @@ -86,7 +90,7 @@ type consumerTask struct { signal string } -func (p *ConsumerPool) AddWork(ctx context.Context, repo string, val *events.XRPCStreamEvent) error { +func (p *Scheduler) AddWork(ctx context.Context, repo string, val *events.XRPCStreamEvent) error { p.itemsAdded.Inc() p.throughputManager.Add(1) t := &consumerTask{ @@ -113,15 +117,15 @@ func (p *ConsumerPool) AddWork(ctx context.Context, repo string, val *events.XRP } } -func (p *ConsumerPool) worker() { +func (p *Scheduler) worker() { log.Infof("starting autoscaling worker for %s", p.ident) - p.workersAcrive.Inc() + p.workersActive.Inc() for work := range p.feeder { for work != nil { // Check if the work item contains a signal to stop the worker. if work.signal == "stop" { log.Infof("stopping autoscaling worker for %s", p.ident) - p.workersAcrive.Dec() + p.workersActive.Dec() return } diff --git a/events/autoscaling/throughput.go b/events/schedulers/autoscaling/throughput.go similarity index 100% rename from events/autoscaling/throughput.go rename to events/schedulers/autoscaling/throughput.go diff --git a/events/schedulers/metrics.go b/events/schedulers/metrics.go new file mode 100644 index 00000000..4b3940ca --- /dev/null +++ b/events/schedulers/metrics.go @@ -0,0 +1,26 @@ +package schedulers + +import ( + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/promauto" +) + +var WorkItemsAdded = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "indigo_scheduler_work_items_added_total", + Help: "Total number of work items added to the consumer pool", +}, []string{"pool", "scheduler_type"}) + +var WorkItemsProcessed = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "indigo_scheduler_work_items_processed_total", + Help: "Total number of work items processed by the consumer pool", +}, []string{"pool", "scheduler_type"}) + +var WorkItemsActive = promauto.NewCounterVec(prometheus.CounterOpts{ + Name: "indigo_scheduler_work_items_active_total", + Help: "Total number of work items passed into a worker", +}, []string{"pool", "scheduler_type"}) + +var WorkersActive = promauto.NewGaugeVec(prometheus.GaugeOpts{ + Name: "indigo_scheduler_workers_active", + Help: "Number of workers currently active", +}, []string{"pool", "scheduler_type"}) diff --git a/events/schedulers/parallel/pool.go b/events/schedulers/parallel/pool.go new file mode 100644 index 00000000..4a0838fb --- /dev/null +++ b/events/schedulers/parallel/pool.go @@ -0,0 +1,117 @@ +package parallel + +import ( + "context" + "sync" + + "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers" + "github.com/labstack/gommon/log" + "github.com/prometheus/client_golang/prometheus" +) + +// Scheduler is a parallel scheduler that will run work on a fixed number of workers +type Scheduler struct { + maxConcurrency int + maxQueue int + + do func(context.Context, *events.XRPCStreamEvent) error + + feeder chan *consumerTask + + lk sync.Mutex + active map[string][]*consumerTask + + ident string + + // metrics + itemsAdded prometheus.Counter + itemsProcessed prometheus.Counter + itemsActive prometheus.Counter + workesActive prometheus.Gauge +} + +func NewScheduler(maxC, maxQ int, ident string, do func(context.Context, *events.XRPCStreamEvent) error) *Scheduler { + p := &Scheduler{ + maxConcurrency: maxC, + maxQueue: maxQ, + + do: do, + + feeder: make(chan *consumerTask), + active: make(map[string][]*consumerTask), + + ident: ident, + + itemsAdded: schedulers.WorkItemsAdded.WithLabelValues(ident, "parallel"), + itemsProcessed: schedulers.WorkItemsProcessed.WithLabelValues(ident, "parallel"), + itemsActive: schedulers.WorkItemsActive.WithLabelValues(ident, "parallel"), + workesActive: schedulers.WorkersActive.WithLabelValues(ident, "parallel"), + } + + for i := 0; i < maxC; i++ { + go p.worker() + } + + p.workesActive.Set(float64(maxC)) + + return p +} + +type consumerTask struct { + repo string + val *events.XRPCStreamEvent +} + +func (p *Scheduler) AddWork(ctx context.Context, repo string, val *events.XRPCStreamEvent) error { + p.itemsAdded.Inc() + t := &consumerTask{ + repo: repo, + val: val, + } + p.lk.Lock() + + a, ok := p.active[repo] + if ok { + p.active[repo] = append(a, t) + p.lk.Unlock() + return nil + } + + p.active[repo] = []*consumerTask{} + p.lk.Unlock() + + select { + case p.feeder <- t: + return nil + case <-ctx.Done(): + return ctx.Err() + } +} + +func (p *Scheduler) worker() { + for work := range p.feeder { + for work != nil { + p.itemsActive.Inc() + if err := p.do(context.TODO(), work.val); err != nil { + log.Errorf("event handler failed: %s", err) + } + p.itemsProcessed.Inc() + + p.lk.Lock() + rem, ok := p.active[work.repo] + if !ok { + log.Errorf("should always have an 'active' entry if a worker is processing a job") + } + + if len(rem) == 0 { + delete(p.active, work.repo) + work = nil + } else { + work = rem[0] + p.active[work.repo] = rem[1:] + } + p.lk.Unlock() + } + } +} diff --git a/events/schedulers/scheduler.go b/events/schedulers/scheduler.go new file mode 100644 index 00000000..9185832f --- /dev/null +++ b/events/schedulers/scheduler.go @@ -0,0 +1 @@ +package schedulers diff --git a/events/schedulers/sequential/sequential.go b/events/schedulers/sequential/sequential.go new file mode 100644 index 00000000..809f32eb --- /dev/null +++ b/events/schedulers/sequential/sequential.go @@ -0,0 +1,47 @@ +package sequential + +import ( + "context" + + "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers" + "github.com/prometheus/client_golang/prometheus" +) + +// Scheduler is a sequential scheduler that will run work on a single worker +type Scheduler struct { + Do func(context.Context, *events.XRPCStreamEvent) error + + ident string + + // metrics + itemsAdded prometheus.Counter + itemsProcessed prometheus.Counter + itemsActive prometheus.Counter + workersActive prometheus.Gauge +} + +func NewScheduler(ident string, do func(context.Context, *events.XRPCStreamEvent) error) *Scheduler { + p := &Scheduler{ + Do: do, + + ident: ident, + + itemsAdded: schedulers.WorkItemsAdded.WithLabelValues(ident, "sequential"), + itemsProcessed: schedulers.WorkItemsProcessed.WithLabelValues(ident, "sequential"), + itemsActive: schedulers.WorkItemsActive.WithLabelValues(ident, "sequential"), + workersActive: schedulers.WorkersActive.WithLabelValues(ident, "sequential"), + } + + p.workersActive.Set(1) + + return p +} + +func (s *Scheduler) AddWork(ctx context.Context, repo string, val *events.XRPCStreamEvent) error { + s.itemsAdded.Inc() + s.itemsActive.Inc() + err := s.Do(ctx, val) + s.itemsProcessed.Inc() + return err +} diff --git a/search/server.go b/search/server.go index 313ed433..66a2624d 100644 --- a/search/server.go +++ b/search/server.go @@ -9,11 +9,13 @@ import ( "net/http" "strconv" "strings" + "time" api "github.com/bluesky-social/indigo/api" comatproto "github.com/bluesky-social/indigo/api/atproto" bsky "github.com/bluesky-social/indigo/api/bsky" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/autoscaling" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/repo" "github.com/bluesky-social/indigo/repomgr" @@ -195,7 +197,7 @@ func (s *Server) RunIndexer(ctx context.Context) error { }, } - return events.HandleRepoStream(ctx, con, events.NewConsumerPool(8, 32, s.bgshost, rsc.EventHandler)) + return events.HandleRepoStream(ctx, con, autoscaling.NewScheduler(1, 32, time.Second, s.bgshost, rsc.EventHandler)) } func (s *Server) handleOp(ctx context.Context, op repomgr.EventKind, seq int64, path string, did string, rcid *cid.Cid, rec any) error { diff --git a/testing/labelmaker_fakedata_test.go b/testing/labelmaker_fakedata_test.go index d11fe78d..d5eb2ad5 100644 --- a/testing/labelmaker_fakedata_test.go +++ b/testing/labelmaker_fakedata_test.go @@ -13,6 +13,7 @@ import ( label "github.com/bluesky-social/indigo/api/label" "github.com/bluesky-social/indigo/carstore" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/sequential" "github.com/bluesky-social/indigo/labeler" "github.com/bluesky-social/indigo/util" "github.com/bluesky-social/indigo/xrpc" @@ -109,7 +110,8 @@ func labelEvents(t *testing.T, lm *labeler.Server, since int64) *EventStream { return nil }, } - if err := events.HandleRepoStream(ctx, con, &events.SequentialScheduler{rsc.EventHandler}); err != nil { + seqScheduler := sequential.NewScheduler("test", rsc.EventHandler) + if err := events.HandleRepoStream(ctx, con, seqScheduler); err != nil { fmt.Println(err) } }() diff --git a/testing/utils.go b/testing/utils.go index d5c821d4..c74487b5 100644 --- a/testing/utils.go +++ b/testing/utils.go @@ -23,6 +23,7 @@ import ( "github.com/bluesky-social/indigo/bgs" "github.com/bluesky-social/indigo/carstore" "github.com/bluesky-social/indigo/events" + "github.com/bluesky-social/indigo/events/schedulers/sequential" "github.com/bluesky-social/indigo/indexer" lexutil "github.com/bluesky-social/indigo/lex/util" "github.com/bluesky-social/indigo/models" @@ -540,7 +541,8 @@ func (b *TestBGS) Events(t *testing.T, since int64) *EventStream { return nil }, } - if err := events.HandleRepoStream(ctx, con, &events.SequentialScheduler{rsc.EventHandler}); err != nil { + seqScheduler := sequential.NewScheduler("test", rsc.EventHandler) + if err := events.HandleRepoStream(ctx, con, seqScheduler); err != nil { fmt.Println(err) } }()