From db4401d42843d321909f17dee3dba44b6ac77ca8 Mon Sep 17 00:00:00 2001 From: Marcus Pasell <3690498+rickyrombo@users.noreply.github.com> Date: Tue, 22 Sep 2026 12:14:15 -0700 Subject: [PATCH] fix(indexer): wait for the core service before starting the ETL MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The ETL's own chain-ID retry budget is 30 attempts at a flat 2s delay — a fixed ~58s, hardcoded as unexported constants in go-openaudio/pkg/etl, so it can't be raised from here. That is shorter than a cold audiusd start: today's node upgrade restarted both at once, audiusd took ~67s to become ready, and the indexer gave up six seconds early. Gate etlIndexer.Run() behind a readiness poll on the same core endpoint the ETL will use, with a 15 minute budget. By the time Run() executes its own retries, the service is already up and the first attempt succeeds. This is the graceful path, not the safety net. The preceding errgroup fix means a failure here now crashes the process and k8s restarts it, but waiting beats crash-looping through escalating backoff and re-running startup on each pass. If core is still down after 15 minutes, crashing is the right answer and that is what happens. Cancellation is distinguished from timeout so a SIGTERM during startup still returns context.Canceled and main.go shuts down quietly instead of panicking. Timeout and poll interval are parameters rather than constants read inside the function, mirroring the upstream initializeChainID signature, so the tests don't sleep for real. Co-Authored-By: Claude Opus 5 --- indexer/core_ready_test.go | 82 ++++++++++++++++++++++++++++++++++++++ indexer/indexer.go | 46 +++++++++++++++++++++ 2 files changed, 128 insertions(+) create mode 100644 indexer/core_ready_test.go diff --git a/indexer/core_ready_test.go b/indexer/core_ready_test.go new file mode 100644 index 00000000..ab701fb1 --- /dev/null +++ b/indexer/core_ready_test.go @@ -0,0 +1,82 @@ +package indexer + +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "connectrpc.com/connect" + corev1 "github.com/OpenAudio/go-openaudio/pkg/api/core/v1" + corev1connect "github.com/OpenAudio/go-openaudio/pkg/api/core/v1/v1connect" + "github.com/OpenAudio/go-openaudio/pkg/sdk" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.uber.org/zap" +) + +type fakeCoreClient struct { + corev1connect.CoreServiceClient + + mu sync.Mutex + calls int + failFirst int +} + +func (f *fakeCoreClient) GetNodeInfo(context.Context, *connect.Request[corev1.GetNodeInfoRequest]) (*connect.Response[corev1.GetNodeInfoResponse], error) { + f.mu.Lock() + defer f.mu.Unlock() + + f.calls++ + if f.failFirst < 0 || f.calls <= f.failFirst { + return nil, connect.NewError(connect.CodeUnavailable, errors.New("core service not ready")) + } + return connect.NewResponse(&corev1.GetNodeInfoResponse{Chainid: "test-chain"}), nil +} + +func (f *fakeCoreClient) callCount() int { + f.mu.Lock() + defer f.mu.Unlock() + return f.calls +} + +func newTestCoreIndexer(core corev1connect.CoreServiceClient) *CoreIndexer { + return &CoreIndexer{ + logger: zap.NewNop(), + openAudioSDK: &sdk.OpenAudioSDK{Core: core}, + } +} + +func TestAwaitCoreReadyRetriesUntilCoreIsUp(t *testing.T) { + core := &fakeCoreClient{failFirst: 3} + ci := newTestCoreIndexer(core) + + err := ci.awaitCoreReady(context.Background(), 5*time.Second, time.Millisecond) + + require.NoError(t, err) + assert.Equal(t, 4, core.callCount()) +} + +func TestAwaitCoreReadyTimesOut(t *testing.T) { + core := &fakeCoreClient{failFirst: -1} + ci := newTestCoreIndexer(core) + + err := ci.awaitCoreReady(context.Background(), 50*time.Millisecond, time.Millisecond) + + require.Error(t, err) + assert.Contains(t, err.Error(), "core service not ready after") + assert.NotErrorIs(t, err, context.Canceled) +} + +func TestAwaitCoreReadyReturnsCanceledOnShutdown(t *testing.T) { + core := &fakeCoreClient{failFirst: -1} + ci := newTestCoreIndexer(core) + + ctx, cancel := context.WithCancel(context.Background()) + cancel() + + err := ci.awaitCoreReady(ctx, time.Minute, time.Millisecond) + + assert.ErrorIs(t, err, context.Canceled) +} diff --git a/indexer/indexer.go b/indexer/indexer.go index 1fcdcf0f..f9b625ed 100644 --- a/indexer/indexer.go +++ b/indexer/indexer.go @@ -2,6 +2,7 @@ package indexer import ( "context" + "errors" "fmt" "net/http" "strings" @@ -12,6 +13,7 @@ import ( "api.audius.co/jobs" "api.audius.co/logging" "connectrpc.com/connect" + corev1 "github.com/OpenAudio/go-openaudio/pkg/api/core/v1" corev1connect "github.com/OpenAudio/go-openaudio/pkg/api/core/v1/v1connect" etl "github.com/OpenAudio/go-openaudio/pkg/etl" em "github.com/OpenAudio/go-openaudio/pkg/etl/processors/entity_manager" @@ -21,6 +23,12 @@ import ( "golang.org/x/sync/errgroup" ) +const ( + coreReadyTimeout = 15 * time.Minute + coreReadyPollInterval = 2 * time.Second + coreReadyLogEvery = 15 +) + // CoreIndexer runs the OpenAudio ETL indexer plus the dependent api/-side // background jobs (aggregates, parity jobs, etc.). The block-fetching and // entity-manager dispatch loop that previously lived here was a stub that @@ -164,6 +172,9 @@ func (ci *CoreIndexer) Start(ctx context.Context) error { return ci.aggregatesCalculator.Start(gCtx) }) eg.Go(func() error { + if err := ci.awaitCoreReady(gCtx, coreReadyTimeout, coreReadyPollInterval); err != nil { + return err + } ci.logger.Info("Starting ETL indexer") return ci.etlIndexer.Run() }) @@ -171,6 +182,41 @@ func (ci *CoreIndexer) Start(ctx context.Context) error { return eg.Wait() } +func (ci *CoreIndexer) awaitCoreReady(ctx context.Context, timeout, pollInterval time.Duration) error { + ctx, cancel := context.WithTimeout(ctx, timeout) + defer cancel() + + var lastErr error + for attempt := 1; ; attempt++ { + _, err := ci.openAudioSDK.Core.GetNodeInfo(ctx, connect.NewRequest(&corev1.GetNodeInfoRequest{})) + if err == nil { + if attempt > 1 { + ci.logger.Info("core service ready", zap.Int("attempts", attempt)) + } + return nil + } + lastErr = err + + if attempt == 1 || attempt%coreReadyLogEvery == 0 { + ci.logger.Warn("waiting for core service", + zap.Int("attempt", attempt), + zap.Duration("timeout", timeout), + zap.Error(err)) + } + + timer := time.NewTimer(pollInterval) + select { + case <-ctx.Done(): + timer.Stop() + if errors.Is(ctx.Err(), context.Canceled) { + return ctx.Err() + } + return fmt.Errorf("core service not ready after %s: %w", timeout, lastErr) + case <-timer.C: + } + } +} + // startParityJobs schedules the periodic jobs that mirror what the legacy // Python discovery-provider celery beat used to run. Each job's ScheduleEvery // launches its own goroutine and exits when ctx is cancelled, so we don't