From 058a2f8d7e6ddbb9c377bbae24b461b7cbd3a5c3 Mon Sep 17 00:00:00 2001 From: Marcus Pasell <3690498+rickyrombo@users.noreply.github.com> Date: Tue, 22 Sep 2026 12:07:28 -0700 Subject: [PATCH] fix(indexer): crash the process when the ETL indexer fails to start MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit CoreIndexer.Start used a plain errgroup.Group, so an error from etlIndexer.Run() was recorded but never propagated: Wait() blocks until every goroutine returns, and AggregatesCalculator.Start only returns on ctx.Done(). With no context linking the two, a fatal ETL error left the process alive with block indexing permanently dead. This bit prod today. A GKE node upgrade restarted every pod at once; audiusd took ~67s to become ready, longer than the ETL's fixed 30x2s chain-ID retry budget, so Run() returned "initialize chain ID after 30 attempts" before indexBlocks() ever started. main.go panics on that error and k8s would have restarted the pod into a healthy audiusd, but Start never returned. The pod reported Running 1/1 for 8 hours while every other background job logged normally and the indexer sat 25k blocks behind, silently dropping plays, follows and every other on-chain write from the API's view. errgroup.WithContext cancels the sibling on the first error, so Wait() now returns the ETL error and main.go panics as intended. Graceful shutdown is unchanged: a SIGTERM-cancelled parent still surfaces as context.Canceled, which main.go ignores. The parity jobs move to the group context for the same reason — they should stop when the process is on its way down. Co-Authored-By: Claude Opus 5 --- indexer/aggregates_calculator_test.go | 27 +++++++++++++++++++++++++++ indexer/indexer.go | 6 +++--- 2 files changed, 30 insertions(+), 3 deletions(-) create mode 100644 indexer/aggregates_calculator_test.go 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() }