diff --git a/indexer/aggregates_calculator_test.go b/indexer/aggregates_calculator_test.go new file mode 100644 index 00000000..f672fe28 --- /dev/null +++ b/indexer/aggregates_calculator_test.go @@ -0,0 +1,27 @@ +package indexer + +import ( + "context" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "go.uber.org/zap" +) + +func TestAggregatesCalculatorStartReturnsOnCancelledContext(t *testing.T) { + a := &AggregatesCalculator{logger: zap.NewNop()} + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + done := make(chan error, 1) + go func() { done <- a.Start(ctx) }() + + select { + case err := <-done: + assert.ErrorIs(t, err, context.Canceled) + case <-time.After(5 * time.Second): + t.Fatal("Start did not return after context cancellation") + } +} diff --git a/indexer/indexer.go b/indexer/indexer.go index e3a133ad..1fcdcf0f 100644 --- a/indexer/indexer.go +++ b/indexer/indexer.go @@ -159,15 +159,15 @@ func newCoreStreamClient(audiusdURL string) corev1connect.CoreServiceClient { // way Go programs always do, and DB connections drain via pool finalizers on // process exit. Acceptable tradeoff to avoid forking ETL. func (ci *CoreIndexer) Start(ctx context.Context) error { - eg := errgroup.Group{} + eg, gCtx := errgroup.WithContext(ctx) eg.Go(func() error { - return ci.aggregatesCalculator.Start(ctx) + return ci.aggregatesCalculator.Start(gCtx) }) eg.Go(func() error { ci.logger.Info("Starting ETL indexer") return ci.etlIndexer.Run() }) - ci.startParityJobs(ctx) + ci.startParityJobs(gCtx) return eg.Wait() }