From 21d183bb65ab8d7b5382689db8ae383cdb3ecf88 Mon Sep 17 00:00:00 2001 From: Terry Tata Date: Tue, 6 Oct 2026 11:32:41 -0700 Subject: [PATCH 1/3] fix(verifier): defer source reader init and tolerate per-chain start failures - sourcereader.Service.Start no longer reads the DB/chain inline; init retries in the background and Ready() reports until it succeeds. - Coordinator.Start skips a chain whose source reader fails to start, keeps the failure visible in HealthReport, and only fails when no chain started at all. - cursechecker.PollerService runs its initial RMN poll in the background goroutine instead of blocking Start. --- integration/pkg/cursechecker/curse_poller.go | 8 +- .../pkg/cursechecker/curse_poller_test.go | 60 ++++++---- verifier/pkg/coordinator.go | 33 +++-- verifier/pkg/coordinator_start_test.go | 83 +++++++++++++ verifier/pkg/sourcereader/service.go | 47 ++++++-- .../pkg/sourcereader/service_start_test.go | 113 ++++++++++++++++++ 6 files changed, 294 insertions(+), 50 deletions(-) create mode 100644 verifier/pkg/coordinator_start_test.go create mode 100644 verifier/pkg/sourcereader/service_start_test.go diff --git a/integration/pkg/cursechecker/curse_poller.go b/integration/pkg/cursechecker/curse_poller.go index 69be32f49..11048c2d2 100644 --- a/integration/pkg/cursechecker/curse_poller.go +++ b/integration/pkg/cursechecker/curse_poller.go @@ -88,9 +88,6 @@ func NewCurseDetectorService( // Start begins polling RMN Remote contracts for curse updates. func (s *PollerService) Start(ctx context.Context) error { return s.StartOnce("cursechecker.PollerService", func() error { - // Initial poll - s.pollAllChains(ctx) - s.wg.Go(func() { s.pollLoop() }) @@ -148,6 +145,11 @@ func (s *PollerService) pollLoop() { ticker := time.NewTicker(s.pollInterval) defer ticker.Stop() + // Initial poll runs in this goroutine: RPC hangs or failures must not + // block Start. Until a poll succeeds, IsRemoteChainCursed reports + // ErrCurseStateUnknown. + s.pollAllChains(ctx) + for { select { case <-ctx.Done(): diff --git a/integration/pkg/cursechecker/curse_poller_test.go b/integration/pkg/cursechecker/curse_poller_test.go index 2fa670297..d8ba48afb 100644 --- a/integration/pkg/cursechecker/curse_poller_test.go +++ b/integration/pkg/cursechecker/curse_poller_test.go @@ -65,11 +65,13 @@ func TestCurseDetectorService_LaneSpecificCurse(t *testing.T) { require.NoError(t, err) defer svc.Close() - cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) - require.NoError(t, err) - assert.True(t, cursed, "chainA->chainB should be cursed") + // The initial poll runs in the background after Start; wait for it to land. + require.Eventually(t, func() bool { + cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) + return err == nil && cursed + }, 2*time.Second, 5*time.Millisecond, "chainA->chainB should be cursed") - cursed, err = svc.IsRemoteChainCursed(ctx, chainA, chainC) + cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainC) require.NoError(t, err) assert.False(t, cursed, "chainA->chainC should not be cursed") @@ -148,10 +150,13 @@ func TestCurseDetectorService_GlobalCurse(t *testing.T) { require.NoError(t, err) defer svc.Close() - cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) - require.NoError(t, err) - assert.True(t, cursed, "chainA has global curse, chainA->chainB should be cursed") - cursed, err = svc.IsRemoteChainCursed(ctx, chainA, chainC) + // The initial poll runs in the background after Start; wait for it to land. + require.Eventually(t, func() bool { + cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) + return err == nil && cursed + }, 2*time.Second, 5*time.Millisecond, "chainA has global curse, chainA->chainB should be cursed") + + cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainC) require.NoError(t, err) assert.True(t, cursed, "chainA has global curse, chainA->chainC should be cursed") @@ -288,9 +293,11 @@ func TestCurseDetectorService_ReaderErrorHandling(t *testing.T) { assert.True(t, cursed, "chainA should be treated as cursed when state unknown (fail closed)") assert.ErrorIs(t, err, common.ErrCurseStateUnknown) - cursed, err = svc.IsRemoteChainCursed(ctx, chainB, chainA) - require.NoError(t, err) - assert.True(t, cursed, "chainB should report chainA as cursed") + // The initial poll runs in the background after Start; wait for it to land. + require.Eventually(t, func() bool { + cursed, err := svc.IsRemoteChainCursed(ctx, chainB, chainA) + return err == nil && cursed + }, 2*time.Second, 5*time.Millisecond, "chainB should report chainA as cursed") } func TestCurseDetectorService_NilCursedSubjects(t *testing.T) { @@ -323,9 +330,11 @@ func TestCurseDetectorService_NilCursedSubjects(t *testing.T) { require.NoError(t, err) defer svc.Close() - cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) - require.NoError(t, err) - assert.False(t, cursed, "nil cursed subjects should mean no curses") + // The initial poll runs in the background after Start; wait for it to land. + require.Eventually(t, func() bool { + cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) + return err == nil && !cursed + }, 2*time.Second, 5*time.Millisecond, "nil cursed subjects should mean no curses") } // TestCurseDetectorService_RPCTimeout tests that hanging RPC calls timeout and don't block other chains. @@ -449,18 +458,17 @@ func TestCurseDetectorService_AllChainsTimeout(t *testing.T) { require.NoError(t, err) defer svc.Close() - // Wait for multiple poll cycles - time.Sleep(400 * time.Millisecond) - - // Verify that polling continued despite timeouts (callCount should be >= 3) - mu.Lock() - count := callCount - mu.Unlock() - assert.GreaterOrEqual(t, count, 3, "polling should continue despite timeouts") - - cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) - require.NoError(t, err) - assert.True(t, cursed, "chainA should eventually report chainB as cursed") + // Verify that polling continued despite timeouts (callCount should reach >= 3) + require.Eventually(t, func() bool { + mu.Lock() + defer mu.Unlock() + return callCount >= 3 + }, 2*time.Second, 20*time.Millisecond, "polling should continue despite timeouts") + + require.Eventually(t, func() bool { + cursed, err := svc.IsRemoteChainCursed(ctx, chainA, chainB) + return err == nil && cursed + }, 2*time.Second, 20*time.Millisecond, "chainA should eventually report chainB as cursed") } // TestCurseDetectorService_ContextCancellation tests that RPC timeout respects parent context cancellation. diff --git a/verifier/pkg/coordinator.go b/verifier/pkg/coordinator.go index f55245aec..44d7c6d00 100644 --- a/verifier/pkg/coordinator.go +++ b/verifier/pkg/coordinator.go @@ -54,6 +54,9 @@ type Coordinator struct { curseDetector common.CurseCheckerService chainStatusBatcher *chainstatus.Batcher sourceReaderServices map[protocol.ChainSelector]services.Service + // sourceReaderStartErrs records per-chain source readers that failed to + // start so HealthReport keeps them visible after they are skipped. + sourceReaderStartErrs map[protocol.ChainSelector]error taskVerifierProcessor services.Service storageWriterProcessor services.Service heartbeatReporter *heartbeat.Reporter @@ -131,10 +134,11 @@ func NewCoordinatorWithDetector( } lggr = logger.With(lggr, "verifierID", config.VerifierID) vc := &Coordinator{ - lggr: lggr, - verifierID: config.VerifierID, - monitoring: monitoring, - messageRulesSvc: messageRulesSvc, + lggr: lggr, + verifierID: config.VerifierID, + monitoring: monitoring, + messageRulesSvc: messageRulesSvc, + sourceReaderStartErrs: make(map[protocol.ChainSelector]error), } vc.initFn = func(ctx context.Context) error { // Batch the chain status writes. The source readers write a status on every @@ -373,13 +377,21 @@ func (vc *Coordinator) Start(ctx context.Context) error { } } - if vc.sourceReaderServices != nil { - for chainSelector, srs := range vc.sourceReaderServices { - if err := srs.Start(ctx); err != nil { - return fmt.Errorf("failed to start source reader service for chain %s: %w", chainSelector, err) - } + // A failure to start one chain's source reader must not stop the remaining + // chains. Log, record for health reporting, and only fail the coordinator + // when no chain started at all. + configuredChains := len(vc.sourceReaderServices) + for chainSelector, srs := range vc.sourceReaderServices { + if err := srs.Start(ctx); err != nil { + vc.lggr.Errorw("Failed to start source reader service, skipping chain", + "chainSelector", chainSelector, "error", err) + vc.sourceReaderStartErrs[chainSelector] = err + delete(vc.sourceReaderServices, chainSelector) } } + if configuredChains > 0 && len(vc.sourceReaderServices) == 0 { + return fmt.Errorf("failed to start any source reader service across %d chains", configuredChains) + } if vc.heartbeatReporter != nil { if err := vc.heartbeatReporter.Start(ctx); err != nil { @@ -578,6 +590,9 @@ func (vc *Coordinator) HealthReport() map[string]error { maps.Copy(report, srs.HealthReport()) } } + for chainSelector, err := range vc.sourceReaderStartErrs { + report[fmt.Sprintf("%s.SourceReader[%s]", vc.Name(), chainSelector)] = err + } if vc.curseDetector != nil { maps.Copy(report, vc.curseDetector.HealthReport()) } diff --git a/verifier/pkg/coordinator_start_test.go b/verifier/pkg/coordinator_start_test.go new file mode 100644 index 000000000..36bfb7c2a --- /dev/null +++ b/verifier/pkg/coordinator_start_test.go @@ -0,0 +1,83 @@ +package verifier + +import ( + "context" + "errors" + "testing" + + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-ccv/protocol" + "github.com/smartcontractkit/chainlink-common/pkg/logger" + "github.com/smartcontractkit/chainlink-common/pkg/services" +) + +// servicesService aliases the interface used for per-chain source readers. +type servicesService = services.Service + +// fakeSourceReaderService is a minimal services.Service whose Start can be +// made to fail, standing in for a per-chain source reader. +type fakeSourceReaderService struct { + name string + startErr error + started bool +} + +func (f *fakeSourceReaderService) Start(context.Context) error { + if f.startErr != nil { + return f.startErr + } + f.started = true + return nil +} + +func (f *fakeSourceReaderService) Close() error { return nil } +func (f *fakeSourceReaderService) Name() string { return f.name } +func (f *fakeSourceReaderService) Ready() error { return nil } +func (f *fakeSourceReaderService) HealthReport() map[string]error { return map[string]error{f.name: nil} } + +// One chain failing to start must not stop the others, and the failure must +// stay visible in the health report. +func TestCoordinator_Start_SkipsFailedSourceReader(t *testing.T) { + const ( + badChain protocol.ChainSelector = 1 + goodChain protocol.ChainSelector = 2 + ) + bad := &fakeSourceReaderService{name: "bad", startErr: errors.New("RPC down")} + good := &fakeSourceReaderService{name: "good"} + + vc := &Coordinator{ + lggr: logger.Test(t), + verifierID: "test-verifier", + monitoring: &noopMonitoring{}, + sourceReaderServices: map[protocol.ChainSelector]servicesService{badChain: bad, goodChain: good}, + sourceReaderStartErrs: make(map[protocol.ChainSelector]error), + } + + require.NoError(t, vc.Start(t.Context())) + defer func() { require.NoError(t, vc.Close()) }() + + require.True(t, good.started, "healthy chain must start despite the failed one") + require.ErrorContains(t, vc.sourceReaderStartErrs[badChain], "RPC down") + require.NotContains(t, vc.sourceReaderServices, badChain, "failed chain is removed from the active set") + + report := vc.HealthReport() + require.ErrorContains(t, report["verifier.Coordinator[test-verifier].SourceReader[1]"], "RPC down") +} + +// If every chain fails to start there is nothing to coordinate: that remains fatal. +func TestCoordinator_Start_FailsWhenNoSourceReaderStarts(t *testing.T) { + vc := &Coordinator{ + lggr: logger.Test(t), + verifierID: "test-verifier", + monitoring: &noopMonitoring{}, + sourceReaderServices: map[protocol.ChainSelector]servicesService{ + 1: &fakeSourceReaderService{name: "a", startErr: errors.New("RPC down")}, + 2: &fakeSourceReaderService{name: "b", startErr: errors.New("RPC down")}, + }, + sourceReaderStartErrs: make(map[protocol.ChainSelector]error), + } + + err := vc.Start(t.Context()) + require.ErrorContains(t, err, "failed to start any source reader service") +} diff --git a/verifier/pkg/sourcereader/service.go b/verifier/pkg/sourcereader/service.go index 33c11511a..b1ec3820f 100644 --- a/verifier/pkg/sourcereader/service.go +++ b/verifier/pkg/sourcereader/service.go @@ -70,7 +70,7 @@ type Service struct { // mutable per-chain state mu sync.RWMutex lastProcessedFinalizedBlock atomic.Uint64 - startBlockInitialized atomic.Bool // guards lastProcessedFinalizedBlock; unset only in unit tests that skip Start + startBlockInitialized atomic.Bool // guards lastProcessedFinalizedBlock; set by the background init in Start pendingTasks map[string]verifier.VerificationTask pendingSince map[string]time.Time pendingMetricDestinations map[protocol.ChainSelector]struct{} @@ -188,18 +188,8 @@ func (r *Service) Start(ctx context.Context) error { return r.StartOnce(r.Name(), func() error { r.logger.Infow("Starting Service") - startBlock, err := r.initializeStartBlock(ctx) - if err != nil { - r.logger.Errorw("Failed to initialize start block", "error", err) - return err - } - r.lastProcessedFinalizedBlock.Store(startBlock) - r.startBlockInitialized.Store(true) - r.metrics().SetSourceReaderLastProcessedFinalizedBlock(ctx, int64(startBlock)) // #nosec G115 -- chain block heights are within int64 range - r.logger.Infow("Initialized start block", "block", startBlock) - r.wg.Go(func() { - r.eventMonitoringLoop() + r.initializeAndMonitor() }) r.logger.Infow("Service started") @@ -207,6 +197,36 @@ func (r *Service) Start(ctx context.Context) error { }) } +// initializeAndMonitor retries start-block initialization until it succeeds or the +// service stops, then hands off to the event monitoring loop. Start must not do +// this inline: the DB/RPC reads can hang or fail on transient errors and must +// never block or abort process startup (Ready reports the interim state instead). +func (r *Service) initializeAndMonitor() { + ctx, cancel := r.stopCh.NewCtx() + defer cancel() + + for { + attemptCtx, attemptCancel := context.WithTimeout(ctx, r.pollTimeout) + startBlock, err := r.initializeStartBlock(attemptCtx) + attemptCancel() + if err == nil { + r.lastProcessedFinalizedBlock.Store(startBlock) + r.startBlockInitialized.Store(true) + r.metrics().SetSourceReaderLastProcessedFinalizedBlock(ctx, int64(startBlock)) // #nosec G115 -- chain block heights are within int64 range + r.logger.Infow("Initialized start block", "block", startBlock) + break + } + r.logger.Errorw("Failed to initialize start block, will retry", "error", err, "retryIn", r.pollInterval) + select { + case <-ctx.Done(): + return + case <-time.After(r.pollInterval): + } + } + + r.eventMonitoringLoop() +} + func (r *Service) Close() error { return r.StopOnce(r.Name(), func() error { r.logger.Infow("Stopping Service") @@ -231,6 +251,9 @@ func (r *Service) Ready() error { if err := r.StateMachine.Ready(); err != nil { return err } + if !r.startBlockInitialized.Load() { + return errors.New("start block not yet initialized") + } if r.finalityBlocked.Load() { return errors.New("finality blocked") } diff --git a/verifier/pkg/sourcereader/service_start_test.go b/verifier/pkg/sourcereader/service_start_test.go new file mode 100644 index 000000000..1821724fc --- /dev/null +++ b/verifier/pkg/sourcereader/service_start_test.go @@ -0,0 +1,113 @@ +package sourcereader + +import ( + "context" + "errors" + "math/big" + "sync/atomic" + "testing" + "time" + + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-ccv/internal/mocks" + "github.com/smartcontractkit/chainlink-ccv/protocol" +) + +// Start must return promptly without waiting on the DB/chain reads that derive +// the start block: those belong to the background loop so a hanging or failing +// dependency cannot block or abort process startup. +func TestSRS_StartDoesNotBlockOnInitIO(t *testing.T) { + chain := protocol.ChainSelector(1337) + + reader := mocks.NewMockSourceReader(t) + chainStatusMgr := mocks.NewMockChainStatusManager(t) + curseDetector := mocks.NewMockCurseCheckerService(t) + + // Hold the DB read open until the test releases it. + readStarted := make(chan struct{}) + release := make(chan struct{}) + chainStatusMgr.EXPECT().ReadChainStatuses(mock.Anything, mock.Anything). + Run(func(context.Context, []protocol.ChainSelector) { + close(readStarted) + <-release + }). + Return(map[protocol.ChainSelector]*protocol.ChainStatusInfo{ + chain: {ChainSelector: chain, FinalizedBlockHeight: big.NewInt(100)}, + }, nil). + Maybe() + reader.EXPECT().LatestAndFinalizedBlock(mock.Anything). + Return(&protocol.BlockHeader{Number: 200}, &protocol.BlockHeader{Number: 150}, nil). + Maybe() + reader.EXPECT().LatestSafeBlock(mock.Anything).Return(nil, nil).Maybe() + reader.EXPECT().FetchMessageSentEvents(mock.Anything, mock.Anything, mock.Anything). + Return([]protocol.MessageSentEvent{}, nil). + Maybe() + curseDetector.EXPECT().IsRemoteChainCursed(mock.Anything, mock.Anything, mock.Anything). + Return(false, nil). + Maybe() + + srs, _, _ := newTestSRS(t, chain, reader, chainStatusMgr, curseDetector, 20*time.Millisecond, 100) + + // Start returns immediately even though the init DB read never completes. + require.NoError(t, srs.Start(t.Context())) + defer func() { require.NoError(t, srs.Close()) }() + + // The background init reaches the blocked read while the service reports not-ready. + select { + case <-readStarted: + case <-time.After(2 * time.Second): + t.Fatal("background init never attempted the chain status read") + } + require.Error(t, srs.Ready(), "service must report not-ready until the start block is initialized") + + // Once the dependency recovers, init completes and the service becomes ready. + close(release) + require.Eventually(t, func() bool { + return srs.Ready() == nil + }, 2*time.Second, 5*time.Millisecond) +} + +// A failing init (DB down, RPC down) must not fail Start; the service retries +// in the background and self-heals. +func TestSRS_StartRetriesFailedInit(t *testing.T) { + chain := protocol.ChainSelector(1337) + + reader := mocks.NewMockSourceReader(t) + chainStatusMgr := mocks.NewMockChainStatusManager(t) + curseDetector := mocks.NewMockCurseCheckerService(t) + + var attempts atomic.Int32 + chainStatusMgr.EXPECT().ReadChainStatuses(mock.Anything, mock.Anything). + RunAndReturn(func(context.Context, []protocol.ChainSelector) (map[protocol.ChainSelector]*protocol.ChainStatusInfo, error) { + if attempts.Add(1) < 3 { + return nil, errors.New("transient DB error") + } + return map[protocol.ChainSelector]*protocol.ChainStatusInfo{ + chain: {ChainSelector: chain, FinalizedBlockHeight: big.NewInt(100)}, + }, nil + }). + Maybe() + reader.EXPECT().LatestAndFinalizedBlock(mock.Anything). + Return(&protocol.BlockHeader{Number: 200}, &protocol.BlockHeader{Number: 150}, nil). + Maybe() + reader.EXPECT().LatestSafeBlock(mock.Anything).Return(nil, nil).Maybe() + reader.EXPECT().FetchMessageSentEvents(mock.Anything, mock.Anything, mock.Anything). + Return([]protocol.MessageSentEvent{}, nil). + Maybe() + curseDetector.EXPECT().IsRemoteChainCursed(mock.Anything, mock.Anything, mock.Anything). + Return(false, nil). + Maybe() + + srs, _, _ := newTestSRS(t, chain, reader, chainStatusMgr, curseDetector, 20*time.Millisecond, 100) + + require.NoError(t, srs.Start(t.Context())) + defer func() { require.NoError(t, srs.Close()) }() + require.Error(t, srs.Ready()) + + require.Eventually(t, func() bool { + return srs.Ready() == nil + }, 2*time.Second, 5*time.Millisecond, "service should self-heal once init succeeds") + require.GreaterOrEqual(t, attempts.Load(), int32(3)) +} From 4e6c2d7ea149d69ac93d3ebb0f331201ba116f68 Mon Sep 17 00:00:00 2001 From: Terry Tata Date: Tue, 6 Oct 2026 12:14:00 -0700 Subject: [PATCH 2/3] fix: tolerate per-chain/per-source startup failures; add noeagerio analyzer Remaining eager-I/O-at-startup fixes beyond the verifier: - coordinator filterConfiguredSourceReaders degrades to unknown statuses instead of failing startup on a transient DB error. - cursechecker runs its initial RMN poll in the background. - token verifier factory skips a failing verifier (fails only if none start); unknown verifier type is a returned config error, not Fatalw. - indexer main skips a failing verifier reader or discovery source (fails only if none start). - pricer skips a chain that fails to start and surfaces it via the new HealthReport; fails only if no chain starts. - aggregator NewServer returns errors instead of Fatalf (signature change: (*Server, error)); main owns the fail-fast decision. Regression prevention: - tools/noeagerio: go/analysis linter flagging I/O (RPC/DB/HTTP/keystore) in New* constructors and Start methods, with intra-package taint propagation and //nolint:noeagerio as the documented escape hatch. Deliberate fail-fast sites (bootstrap DB/keystore, JD job load, signer key load, aggregator storage, replay tool) are annotated. - Wired into just lint-noeagerio and the golangci-lint CI workflow. - Policy recorded in AGENTS.md. --- .github/workflows/golangci-lint.yaml | 4 + AGENTS.md | 13 + Justfile | 6 + aggregator/cmd/main.go | 5 +- .../rate_limiter_store_factory.go | 1 + aggregator/pkg/server.go | 16 +- aggregator/tests/utils.go | 6 +- bootstrap/bootstrap.go | 2 + bootstrap/keys/csa.go | 1 + bootstrap/keys/kms.go | 1 + changelog/2026-10-06_no_eager_io_startup.md | 35 ++ cmd/verifier/tokenfactory.go | 33 +- common/jd/lifecycle/manager.go | 1 + go.mod | 2 +- indexer/cmd/main.go | 24 +- indexer/pkg/replay/engine.go | 2 + indexer/pkg/replay/store.go | 1 + indexer/pkg/storage/postgres.go | 1 + .../evm_contract_transmitter.go | 1 + pricer/pkg/coordinator/coordinator.go | 41 +- pricer/pkg/evm/evm.go | 5 +- pricer/pkg/sol/sol.go | 4 + tools/noeagerio/README.md | 39 ++ tools/noeagerio/analyzer.go | 418 ++++++++++++++++++ tools/noeagerio/analyzer_test.go | 11 + tools/noeagerio/cmd/noeagerio/main.go | 17 + tools/noeagerio/testdata/src/a/a.go | 99 +++++ .../testdata/src/common/lazy/lazy.go | 17 + verifier/pkg/commit/signer.go | 1 + verifier/pkg/coordinator.go | 25 +- verifier/pkg/coordinator_start_test.go | 10 +- 31 files changed, 794 insertions(+), 48 deletions(-) create mode 100644 changelog/2026-10-06_no_eager_io_startup.md create mode 100644 tools/noeagerio/README.md create mode 100644 tools/noeagerio/analyzer.go create mode 100644 tools/noeagerio/analyzer_test.go create mode 100644 tools/noeagerio/cmd/noeagerio/main.go create mode 100644 tools/noeagerio/testdata/src/a/a.go create mode 100644 tools/noeagerio/testdata/src/common/lazy/lazy.go diff --git a/.github/workflows/golangci-lint.yaml b/.github/workflows/golangci-lint.yaml index f9b29a572..4cc5037f9 100644 --- a/.github/workflows/golangci-lint.yaml +++ b/.github/workflows/golangci-lint.yaml @@ -63,3 +63,7 @@ jobs: - name: Lint run: just lint --timeout 15m + + # Custom analyzer: no I/O in constructors or Start methods (tools/noeagerio). + - name: Lint (noeagerio) + run: just lint-noeagerio diff --git a/AGENTS.md b/AGENTS.md index 6a068fb00..96c812436 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -25,3 +25,16 @@ stored in JD. Reference where the secret lives; keep the value in the secrets fi Long comment blocks are hard to comprehend and takes a lot of space. State the rule/behavior in one line, the reason in one more if it's non-obvious, and stop. If a comment needs a paragraph, put that reasoning in the PR description or commit message instead. + +## Startup: no eager I/O in constructors or Start() + +Eager I/O at startup has caused multiple production outages (a rate-limited RPC +took down the whole committee verifier). Constructors (`New*`) only assemble +dependencies; `Start(ctx)` only spawns goroutines and returns. Any network, DB, +RPC, or keystore read must happen at query time — via `common/lazy.Lazy` for +cached-on-success derivations — or in a background goroutine that reports its +state through `Ready()`/`HealthReport()`. Per-chain and per-source failures +degrade: log, skip, record in the health report, and fail startup only when +*nothing* usable remains. The `noeagerio` analyzer (`tools/noeagerio`, run by +`just lint-noeagerio` and CI) enforces this; a deliberate fail-fast exception +(identity keys, the service's own DB) needs `//nolint:noeagerio` with a reason. diff --git a/Justfile b/Justfile index a33c92608..05bb1cad7 100644 --- a/Justfile +++ b/Justfile @@ -66,6 +66,12 @@ fmt: ensure-golangci-lint lint fix="" timeout="10m": ensure-golangci-lint gomods -c 'golangci-lint run --config {{justfile_directory()}}/.golangci.yaml {{ if fix == "true" { "--fix" } else if fix == "fix" { "--fix" } else { "" } }} --timeout {{timeout}}' +# Run the noeagerio analyzer (tools/noeagerio): no I/O in constructors or Start +# methods. Runs on the root module only; build/devenv and deployment are excluded. +lint-noeagerio: ensure-go + go build -o "$(go env GOPATH)/bin/noeagerio" ./tools/noeagerio/cmd/noeagerio + go vet -vettool="$(go env GOPATH)/bin/noeagerio" ./... + shellcheck: @command -v shellcheck >/dev/null 2>&1 || { \ echo "shellcheck is not installed. Please install it first."; \ diff --git a/aggregator/cmd/main.go b/aggregator/cmd/main.go index 1fe717870..7fd9d0b65 100644 --- a/aggregator/cmd/main.go +++ b/aggregator/cmd/main.go @@ -210,7 +210,10 @@ func runServer(configPath, logLevelStr string, lggr logger.Logger, sugaredLggr l protocol.InitChainSelectorCache() - server := aggregator.NewServer(sugaredLggr, config, aggMonitoring) + server, err := aggregator.NewServer(sugaredLggr, config, aggMonitoring) + if err != nil { + sugaredLggr.Fatalw("failed to create CCV data service", "error", err) + } ctx := context.Background() ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) defer stop() diff --git a/aggregator/pkg/rate_limiting/rate_limiter_store_factory.go b/aggregator/pkg/rate_limiting/rate_limiter_store_factory.go index 412a5ad88..eb6b86aed 100644 --- a/aggregator/pkg/rate_limiting/rate_limiter_store_factory.go +++ b/aggregator/pkg/rate_limiting/rate_limiter_store_factory.go @@ -35,6 +35,7 @@ func NewRateLimiterStore(config model.RateLimiterStoreConfig) (limiter.Store, er DB: config.Redis.DB, }) + //nolint:noeagerio // fail-fast by design: rate limiting is opt-in protection, so starting up with an unreachable store beats silently running unprotected (Ready() also health-checks the store) if err := redisClient.Ping(context.Background()).Err(); err != nil { return nil, fmt.Errorf("failed to connect to redis at %s: %w", config.Redis.Address, err) } diff --git a/aggregator/pkg/server.go b/aggregator/pkg/server.go index ab948ca2f..f9992f8ce 100644 --- a/aggregator/pkg/server.go +++ b/aggregator/pkg/server.go @@ -309,9 +309,11 @@ type SignatureAndQuorumValidator interface { // NewServer creates a new aggregator server with the specified logger, configuration, and monitoring. // aggMonitoring must not be nil; use monitoring.NoopAggregatorMonitoring when monitoring is disabled. -func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonitoring common.AggregatorMonitoring) *Server { +// Errors are returned to the caller (main), which owns the fail-fast decision; a +// library constructor must never exit the process. +func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonitoring common.AggregatorMonitoring) (*Server, error) { if err := config.Validate(); err != nil { - l.Fatalf("Failed to validate server configuration: %v", err) + return nil, fmt.Errorf("failed to validate server configuration: %w", err) } l.Infow("Server configuration loaded", @@ -325,10 +327,10 @@ func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonito ) factory := storage.NewStorageFactory(l) + //nolint:noeagerio // the aggregator's own DB is a hard dependency: connect + migrate fail fast at startup, and the health endpoint reports readiness after that rawStore, err := factory.CreateStorage(config.Storage, aggMonitoring) if err != nil { - l.Fatalf("Failed to create storage: %v", err) - return nil + return nil, fmt.Errorf("failed to create storage: %w", err) } // Build the message-disablement registry from the raw store before metrics wrapping. @@ -365,14 +367,14 @@ func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonito hmacAuthMiddleware := middlewares.NewHMACAuthMiddleware(config, l) anonymousAuthMiddleware, err := middlewares.NewAnonymousAuthMiddleware(config.AnonymousAuth.TrustedProxies, l) if err != nil { - l.Fatalf("Failed to initialize anonymous auth middleware: %v", err) + return nil, fmt.Errorf("failed to initialize anonymous auth middleware: %w", err) } requireAuthMiddleware := middlewares.NewRequireAuthMiddleware(l) // Initialize rate limiting middleware rateLimitingMiddleware, err := middlewares.NewRateLimitingMiddlewareFromConfig(config.RateLimiting, config, l) if err != nil { - l.Fatalf("Failed to initialize rate limiting middleware: %v", err) + return nil, fmt.Errorf("failed to initialize rate limiting middleware: %w", err) } isVerifierResultAPI := func(callMeta interceptors.CallMeta) bool { @@ -470,5 +472,5 @@ func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonito committeepb.RegisterCommitteeVerifierServer(grpcServer, server) heartbeatpb.RegisterHeartbeatServiceServer(grpcServer, server) - return server + return server, nil } diff --git a/aggregator/tests/utils.go b/aggregator/tests/utils.go index 600478547..f1ed8f8c7 100644 --- a/aggregator/tests/utils.go +++ b/aggregator/tests/utils.go @@ -233,7 +233,11 @@ func CreateServerOnlyWithMessageRulesControl(t *testing.T, options ...ConfigOpti return nil, nil, nil, fmt.Errorf("test storage does not implement message rules store") } - s := agg.NewServer(sugaredLggr, config, monitoring.NewNoopAggregatorMonitoring()) + s, err := agg.NewServer(sugaredLggr, config, monitoring.NewNoopAggregatorMonitoring()) + if err != nil { + cleanupStorage() + return nil, nil, nil, fmt.Errorf("failed to create server: %w", err) + } err = s.Start(buf) if err != nil { t.Fatalf("failed to start server: %v", err) diff --git a/bootstrap/bootstrap.go b/bootstrap/bootstrap.go index 3d7ecfc78..45d1275bb 100644 --- a/bootstrap/bootstrap.go +++ b/bootstrap/bootstrap.go @@ -710,8 +710,10 @@ func (b *Bootstrapper) Start(ctx context.Context) error { return fmt.Errorf("bootstrapper has no logger") } if b.mode == AppConfigModeJD { + //nolint:noeagerio // fail-fast by design: the DB connection, migrations, and keystore are hard dependencies of any job, verified up front and bounded by the startup timeout return b.startWithJDLifecycle(ctx) } + //nolint:noeagerio // fail-fast by design: same as the JD path above; a misconfigured local deployment should exit, not idle return b.startLocal(ctx) } diff --git a/bootstrap/keys/csa.go b/bootstrap/keys/csa.go index ce77da7aa..c53c33308 100644 --- a/bootstrap/keys/csa.go +++ b/bootstrap/keys/csa.go @@ -38,6 +38,7 @@ var _ crypto.Signer = (*CSASigner)(nil) // NewCSASigner returns a [crypto.Signer] for the named Ed25519 key in ks. func NewCSASigner(ctx context.Context, ks keystore.Keystore, keyName string) (*CSASigner, error) { + //nolint:noeagerio // fail-fast by design: without the CSA key the node cannot identify itself, so the keystore read happens at startup rather than at first use resp, err := ks.GetKeys(ctx, keystore.GetKeysRequest{ KeyNames: []string{keyName}, }) diff --git a/bootstrap/keys/kms.go b/bootstrap/keys/kms.go index 85243af14..56b63b6ab 100644 --- a/bootstrap/keys/kms.go +++ b/bootstrap/keys/kms.go @@ -88,6 +88,7 @@ func newKMSKeystore(ctx context.Context, inner signerReader, nameToID map[string // verifyKeys checks that every mapped KMS key exists and is accessible. func (k *KMSKeystore) verifyKeys(ctx context.Context) error { for name, id := range k.nameToID { + //nolint:noeagerio // fail-fast by design: the service can sign nothing if a mapped KMS key is missing, so key existence is verified at startup rather than surfacing at first use _, err := k.inner.GetKeys(ctx, keystore.GetKeysRequest{KeyNames: []string{id}}) if err != nil { return fmt.Errorf("KMS key %q (logical name %q) not accessible: %w", id, name, err) diff --git a/changelog/2026-10-06_no_eager_io_startup.md b/changelog/2026-10-06_no_eager_io_startup.md new file mode 100644 index 000000000..efb9804a5 --- /dev/null +++ b/changelog/2026-10-06_no_eager_io_startup.md @@ -0,0 +1,35 @@ +# Human Overview + +Reliability hardening against eager I/O at startup, the failure class behind +repeated verifier outages (see also PR #1424): + +* `verifier/pkg/sourcereader`: `Start` no longer reads the DB/chain inline; + start-block initialization retries in the background and `Ready()` reports + until it succeeds. A transient RPC/DB failure can no longer abort startup. +* `verifier/pkg` coordinator: a per-chain source reader that fails to start is + skipped, recorded in `HealthReport()`, and only fails the coordinator when + no chain started at all. The all-chains chain-status read at coordinator + start degrades to "unknown" instead of failing startup. +* `integration/pkg/cursechecker`: the initial RMN poll runs in the background + goroutine instead of blocking `Start`. +* `cmd/verifier` token factory: one failing token verifier no longer prevents + the others from starting (and `Fatalw` on unknown verifier type is now a + returned error). Fails only when no verifier started. +* `indexer/cmd`: one failing verifier reader or discovery source no longer + exits the process; fails only when none started. +* `pricer`: a chain that fails to start is skipped and surfaced via the new + `HealthReport()`; fails only when no chain started. +* `aggregator/pkg`: `NewServer` returns errors instead of calling + `logger.Fatalf` (signature changed to `(*Server, error)`); the caller in + `main` owns the fail-fast decision. +* New `noeagerio` static analyzer (`tools/noeagerio`, run by + `just lint-noeagerio` and the lint CI workflow) forbids I/O in constructors + and `Start` methods repo-wide, with `//nolint:noeagerio` as the documented + escape hatch for deliberate fail-fast exceptions. Policy recorded in + AGENTS.md. + +Deliberately unchanged (fail-fast by design, annotated with `//nolint:noeagerio` ++ justification): bootstrap DB connect/migrations and keystore/KMS key +verification, the JD lifecycle cached-job load, the commit signer keystore +read, the aggregator's own storage connect+migrate, the opt-in Redis rate +limiter, and the operator-run indexer replay tool. diff --git a/cmd/verifier/tokenfactory.go b/cmd/verifier/tokenfactory.go index 1acbdb210..1e71f8044 100644 --- a/cmd/verifier/tokenfactory.go +++ b/cmd/verifier/tokenfactory.go @@ -150,9 +150,10 @@ func (tvf *tokenVerifierFactory) Start(ctx context.Context, spec bootstrap.JobSp ) var coordinator *verifier.Coordinator + var createErr error switch { case verifierConfig.IsLombard(): - coordinator, err = createLombardCoordinator( + coordinator, createErr = createLombardCoordinator( ctx, verifierConfig.VerifierID, verifierConfig.LombardConfig, @@ -172,9 +173,10 @@ func (tvf *tokenVerifierFactory) Start(ctx context.Context, spec bootstrap.JobSp case verifierConfig.IsCCTP(): cctpCodecs, codecErr := cctpCodecsFor(accessors, verifierConfig.CCTPConfig) if codecErr != nil { - return fmt.Errorf("failed to create verification coordinator for cctp: %w", codecErr) + tvf.lggr.Errorw("Skipping verifier, failed to build CCTP codecs", "verifierID", verifierConfig.VerifierID, "error", codecErr) + continue } - coordinator, err = createCCTPCoordinator( + coordinator, createErr = createCCTPCoordinator( ctx, verifierConfig.VerifierID, verifierConfig.CCTPConfig, @@ -193,19 +195,28 @@ func (tvf *tokenVerifierFactory) Start(ctx context.Context, spec bootstrap.JobSp db, ) default: - tvf.lggr.Fatalw("Unknown verifier type", "type", verifierConfig.Type) - continue + // Unknown type is a deterministic config error, not a transient + // failure: fail fast rather than silently skipping the verifier. + return fmt.Errorf("unknown verifier type %q for verifier %s", verifierConfig.Type, verifierConfig.VerifierID) } - if err != nil { - return fmt.Errorf("failed to create verification coordinator for %s: %w", verifierConfig.Type, err) + // A failure to stand up one verifier must not stop the remaining + // verifiers. Log and skip; only reject the whole service when no + // verifier is usable. + if createErr != nil { + tvf.lggr.Errorw("Skipping verifier, failed to create verification coordinator", + "verifierID", verifierConfig.VerifierID, "type", verifierConfig.Type, "error", createErr) + continue } - tvf.coordinators = append(tvf.coordinators, coordinator) - if err := coordinator.Start(ctx); err != nil { - tvf.lggr.Errorw("Failed to start verification coordinator", "error", err) - return fmt.Errorf("failed to start verification coordinator: %w", err) + tvf.lggr.Errorw("Skipping verifier, failed to start verification coordinator", + "verifierID", verifierConfig.VerifierID, "type", verifierConfig.Type, "error", err) + continue } + tvf.coordinators = append(tvf.coordinators, coordinator) + } + if len(cfg.TokenVerifiers) > 0 && len(tvf.coordinators) == 0 { + return fmt.Errorf("failed to start any token verifier across %d configured", len(cfg.TokenVerifiers)) } healthReporters := make([]protocol.HealthReporter, len(tvf.coordinators)) diff --git a/common/jd/lifecycle/manager.go b/common/jd/lifecycle/manager.go index 0b7cd6354..56caa4eef 100644 --- a/common/jd/lifecycle/manager.go +++ b/common/jd/lifecycle/manager.go @@ -165,6 +165,7 @@ func (m *Manager) Start(ctx context.Context) error { m.lggr.Infow("Starting job lifecycle manager") // 1. Load cached job if exists + //nolint:noeagerio // fail-fast by design: the cached job determines what this process runs; proceeding without it would start the wrong (or no) job cachedJob, err := m.jobStore.LoadJob(ctx) if err != nil && !errors.Is(err, store.ErrNoJob) { return fmt.Errorf("failed to load cached job: %w", err) diff --git a/go.mod b/go.mod index 97d39fd64..904b99951 100644 --- a/go.mod +++ b/go.mod @@ -62,6 +62,7 @@ require ( golang.org/x/mod v0.38.0 golang.org/x/sync v0.22.0 golang.org/x/time v0.15.0 + golang.org/x/tools v0.48.0 google.golang.org/genproto/googleapis/rpc v0.0.0-20260819154853-08b0e4226688 google.golang.org/grpc v1.83.2 google.golang.org/protobuf v1.36.12 @@ -341,7 +342,6 @@ require ( golang.org/x/sys v0.47.0 // indirect golang.org/x/term v0.45.0 // indirect golang.org/x/text v0.41.0 // indirect - golang.org/x/tools v0.48.0 // indirect google.golang.org/api v0.287.1 // indirect google.golang.org/genproto v0.0.0-20260319201613-d00831a3d3e7 // indirect google.golang.org/genproto/googleapis/api v0.0.0-20260819154853-08b0e4226688 // indirect diff --git a/indexer/cmd/main.go b/indexer/cmd/main.go index 32786f75c..3cb0c45ff 100644 --- a/indexer/cmd/main.go +++ b/indexer/cmd/main.go @@ -158,11 +158,20 @@ func createRegistry() *registry.VerifierRegistry { } func createAllVerifierReaders(ctx context.Context, lggr logger.Logger, verifierRegistry *registry.VerifierRegistry, config *config.Config, indexerMonitoring common.IndexerMonitoring) error { + // A failure to stand up one verifier's reader (e.g. a bad address) must not + // stop the remaining verifiers. Log and skip; only fail when no verifier + // is usable. + created := 0 for _, verifierConfig := range config.Verifiers { err := createReadersForVerifier(ctx, lggr, verifierRegistry, &verifierConfig, indexerMonitoring, config.Resilience) if err != nil { - return err + lggr.Errorw("Skipping verifier, failed to create readers", "verifier", verifierConfig.Name, "error", err) + continue } + created++ + } + if len(config.Verifiers) > 0 && created == 0 { + return fmt.Errorf("failed to create readers for any of the %d configured verifiers", len(config.Verifiers)) } return nil @@ -252,8 +261,9 @@ func createDiscovery(ctx context.Context, lggr logger.Logger, cfg *config.Config Secret: discCfg.Secret, }, discCfg.InsecureConnection, config.EffectiveMaxResponseBytes(discCfg.MaxResponseBytes), metrics, readers.NewResilienceConfig(cfg.Resilience)) if err != nil { - cleanupOnError() - return nil, err + // One misconfigured discovery source must not stop the others. + lggr.Errorw("Skipping discovery source, failed to create aggregator reader", "address", discCfg.Address, "error", err) + continue } ntpKey := fmt.Sprintf("%s|%d", discCfg.NtpServer, discCfg.Timeout) @@ -276,12 +286,16 @@ func createDiscovery(ctx context.Context, lggr logger.Logger, cfg *config.Config discovery.WithPrimaryWriteNotifier(writeNotifier), // nil for single-source; no-op ) if err != nil { - cleanupOnError() - return nil, err + // One misconfigured discovery source must not stop the others. + lggr.Errorw("Skipping discovery source, failed to create message discovery", "address", discCfg.Address, "error", err) + continue } sources = append(sources, aggDiscovery) } + if len(configs) > 0 && len(sources) == 0 { + return nil, fmt.Errorf("failed to create any discovery source across %d configured", len(configs)) + } if len(sources) == 1 { return sources[0], nil } diff --git a/indexer/pkg/replay/engine.go b/indexer/pkg/replay/engine.go index 75a7e4c80..c3404bfd2 100644 --- a/indexer/pkg/replay/engine.go +++ b/indexer/pkg/replay/engine.go @@ -69,6 +69,8 @@ func NewEngine( // Start creates a new job or resumes a stale one, acquires the advisory lock, // and runs the replay to completion. +// +//nolint:noeagerio // replay is an operator-run one-shot tool: Start IS the job, not a service lifecycle, so DB work here is the program doing its job func (e *Engine) Start(ctx context.Context, req Request) (string, error) { job, err := e.findOrCreateJob(ctx, req) if err != nil { diff --git a/indexer/pkg/replay/store.go b/indexer/pkg/replay/store.go index 86a540336..cfc36da66 100644 --- a/indexer/pkg/replay/store.go +++ b/indexer/pkg/replay/store.go @@ -34,6 +34,7 @@ func NewStore(ds sqlutil.DataSource, lggr logger.Logger) *Store { // NewStoreFromConfig creates a Store with its own Postgres connection pool. func NewStoreFromConfig(ctx context.Context, lggr logger.Logger, uri string, dbConfig pg.DBConfig, connMaxLifetime, connMaxIdleTime time.Duration) (*Store, error) { + //nolint:noeagerio // replay is an operator-run one-shot tool whose only dependency is this DB; connecting eagerly fails fast on misconfiguration db, err := dbConfig.New(ctx, uri, pg.DriverPostgres) if err != nil { return nil, fmt.Errorf("failed to open replay store connection: %w", err) diff --git a/indexer/pkg/storage/postgres.go b/indexer/pkg/storage/postgres.go index 60fc0326f..54104e2fd 100644 --- a/indexer/pkg/storage/postgres.go +++ b/indexer/pkg/storage/postgres.go @@ -70,6 +70,7 @@ func NewPostgresStorage(ctx context.Context, lggr logger.Logger, monitoring comm if driverName == "" { return nil, fmt.Errorf("database driver name is required") } + //nolint:noeagerio // fail-fast by design: the service's own DB is a hard dependency; callers wire this from main where startup failure exits the process db, err := config.New(ctx, uri, driverName) if err != nil { return nil, fmt.Errorf("failed to open database connection: %w", err) diff --git a/integration/pkg/contracttransmitter/evm_contract_transmitter.go b/integration/pkg/contracttransmitter/evm_contract_transmitter.go index 41db0bca6..b3ea60457 100644 --- a/integration/pkg/contracttransmitter/evm_contract_transmitter.go +++ b/integration/pkg/contracttransmitter/evm_contract_transmitter.go @@ -38,6 +38,7 @@ type EVMContractTransmitter struct { func NewEVMContractTransmitterFromRPC(_ context.Context, lggr logger.Logger, chainSelector protocol.ChainSelector, rpc, privatekey string, offRampAddress common.Address) (*EVMContractTransmitter, error) { // create a client for the off ramp contract + //nolint:noeagerio // retained for external callers (CL node integration); a contract transmitter is useless without its chain connection, so dialing fails fast client, err := ethclient.Dial(rpc) if err != nil { return nil, err diff --git a/pricer/pkg/coordinator/coordinator.go b/pricer/pkg/coordinator/coordinator.go index e2e98b6c3..8b18f9d20 100644 --- a/pricer/pkg/coordinator/coordinator.go +++ b/pricer/pkg/coordinator/coordinator.go @@ -118,6 +118,9 @@ type Pricer struct { wg sync.WaitGroup httpServer *http.Server chains map[protocol.ChainSelector]pricer.Chain + // chainStartErrs records per-chain start failures so they stay visible in + // HealthReport after the chain is skipped. + chainStartErrs map[protocol.ChainSelector]error } func NewPricerFromConfig(ctx context.Context, cfg Config, keystoreData []byte, keystorePassword string) (*Pricer, error) { @@ -199,11 +202,12 @@ func NewPricerFromConfig(ctx context.Context, cfg Config, keystoreData []byte, k )) return &Pricer{ - StateMachine: services.StateMachine{}, - lggr: lggr, - cfg: cfg, - done: make(chan struct{}), - wg: sync.WaitGroup{}, + StateMachine: services.StateMachine{}, + lggr: lggr, + cfg: cfg, + done: make(chan struct{}), + wg: sync.WaitGroup{}, + chainStartErrs: make(map[protocol.ChainSelector]error), httpServer: &http.Server{ Addr: fmt.Sprintf(":%d", cfg.Monitoring.Port), Handler: mux, @@ -227,11 +231,20 @@ func (p *Pricer) Start(ctx context.Context) error { } }) - for _, chain := range p.chains { + // A failure to start one chain (e.g. an unreachable RPC) must not stop + // the remaining chains. Skip and record for health reporting; only fail + // when no chain started at all. + configuredChains := len(p.chains) + for selector, chain := range p.chains { if err := chain.Start(ctx); err != nil { - return fmt.Errorf("failed to start chain: %w", err) + p.lggr.Errorw("failed to start chain, skipping", "chainSelector", selector, "error", err) + p.chainStartErrs[selector] = err + delete(p.chains, selector) } } + if configuredChains > 0 && len(p.chains) == 0 { + return fmt.Errorf("failed to start any chain across %d configured", configuredChains) + } p.wg.Go(func() { p.run(ctx) }) @@ -297,3 +310,17 @@ func (p *Pricer) Close() error { return nil }) } + +func (p *Pricer) Name() string { return "pricer.Pricer" } + +// HealthReport surfaces chains that failed to start so a skipped chain is +// visible instead of silently absent. +func (p *Pricer) HealthReport() map[string]error { + report := map[string]error{p.Name(): p.Ready()} + for selector, err := range p.chainStartErrs { + report[fmt.Sprintf("%s.Chain[%s]", p.Name(), selector)] = err + } + return report +} + +var _ protocol.HealthReporter = (*Pricer)(nil) diff --git a/pricer/pkg/evm/evm.go b/pricer/pkg/evm/evm.go index cf70c8f55..02d2e046a 100644 --- a/pricer/pkg/evm/evm.go +++ b/pricer/pkg/evm/evm.go @@ -36,7 +36,10 @@ type Chain struct { } func (c *Chain) Start(ctx context.Context) error { - // Dial the EVM client to start the connection pool. + // Dial the EVM client to start the connection pool. The TXM requires a live + // connection, so this stays eager; the pricer coordinator tolerates a + // per-chain start failure and reports it via HealthReport. + //nolint:noeagerio // see comment above if err := c.client.Dial(ctx); err != nil { return fmt.Errorf("failed to dial EVM client: %w", err) } diff --git a/pricer/pkg/sol/sol.go b/pricer/pkg/sol/sol.go index a9e063a53..5f2cf7b6a 100644 --- a/pricer/pkg/sol/sol.go +++ b/pricer/pkg/sol/sol.go @@ -30,6 +30,10 @@ type Chain struct { } func (c *Chain) Start(ctx context.Context) error { + // The TXM requires a live connection, so the dial stays eager; the pricer + // coordinator tolerates a per-chain start failure and reports it via + // HealthReport. + //nolint:noeagerio // see comment above if err := c.client.Dial(ctx); err != nil { return fmt.Errorf("failed to dial Solana client: %w", err) } diff --git a/tools/noeagerio/README.md b/tools/noeagerio/README.md new file mode 100644 index 000000000..8fbc95119 --- /dev/null +++ b/tools/noeagerio/README.md @@ -0,0 +1,39 @@ +# noeagerio + +`noeagerio` is a `go/analysis` linter that forbids I/O (RPC, DB, HTTP, +keystore/KMS) inside startup paths: exported `New*` constructors and `Start` +methods. See AGENTS.md ("Startup: no eager I/O in constructors or Start()") +for the policy and the outage history behind it. + +Run it with: + +```sh +just lint-noeagerio +``` + +## What it flags + +- Direct calls to a denylist of I/O functions/methods (see `denyFuncs` and + `denyMethods` in `analyzer.go`) anywhere in the synchronous body of a + constructor or `Start` method. +- Calls to in-package helpers that (transitively) perform such I/O, reported + at the call site in the startup function. + +## What it deliberately allows + +- I/O inside goroutines spawned from `Start` (`go func(){...}`, + `wg.Go(func(){...})`, errgroup, ...) — the blessed pattern, provided the + service reports progress via `Ready()`/`HealthReport()`. +- Closures registered for query-time execution (`common/lazy.New`). +- Anything outside constructors/`Start`. +- Call sites carrying `//nolint:noeagerio` with a justification (on the line, + the line above, or the enclosing function's doc comment). A suppressed call + does not taint its callers. + +## Limitations + +Taint propagation is intra-package only; the cross-package cases we know about +are covered by repo-specific denylist entries. A determined author can evade +the check (e.g. hiding I/O behind an unexported helper called from an +interface) — the analyzer is a tripwire, not a proof. Code review remains the +backstop. diff --git a/tools/noeagerio/analyzer.go b/tools/noeagerio/analyzer.go new file mode 100644 index 000000000..591685a11 --- /dev/null +++ b/tools/noeagerio/analyzer.go @@ -0,0 +1,418 @@ +// Package noeagerio implements a static analysis that flags I/O performed +// inside startup paths: constructors (exported New* functions) and Start +// methods. Startup I/O has repeatedly caused outages in this repo — a hanging +// or failing RPC/DB call in a constructor or Start blocks or aborts process +// startup, and a single bad chain can take down every healthy one. +// +// I/O belongs at query time (see common/lazy for cached-on-success derivation) +// or in a background goroutine that reports its state via Ready/HealthReport. +// Deliberate exceptions need a //nolint:noeagerio comment with a justification. +package noeagerio + +import ( + "go/ast" + "go/token" + "go/types" + "strings" + + "golang.org/x/tools/go/analysis" + "golang.org/x/tools/go/analysis/passes/inspect" + "golang.org/x/tools/go/ast/inspector" +) + +// Analyzer is the go/analysis entry point. +var Analyzer = &analysis.Analyzer{ + Name: "noeagerio", + Doc: "reports I/O calls in constructors (New*) and Start methods; defer I/O to query time or background goroutines", + Run: run, + Requires: []*analysis.Analyzer{inspect.Analyzer}, +} + +// denyFuncs are package-level functions that perform network/DB I/O, keyed by +// "import/path.Name". +var denyFuncs = map[string]string{ + "net/http.Get": "issues an HTTP request", + "net/http.Head": "issues an HTTP request", + "net/http.Post": "issues an HTTP request", + "net/http.PostForm": "issues an HTTP request", + "github.com/jmoiron/sqlx.Connect": "opens and pings a database connection", + "github.com/jmoiron/sqlx.ConnectContext": "opens and pings a database connection", + "github.com/jmoiron/sqlx.MustConnect": "opens and pings a database connection", + "google.golang.org/grpc.Dial": "dials a gRPC endpoint", + "google.golang.org/grpc.DialContext": "dials a gRPC endpoint", + "github.com/ethereum/go-ethereum/ethclient.Dial": "dials an Ethereum RPC endpoint", + "github.com/ethereum/go-ethereum/ethclient.DialContext": "dials an Ethereum RPC endpoint", + "github.com/jackc/pgx/v5.Connect": "opens a database connection", + "github.com/jackc/pgx/v5/pgxpool.Connect": "opens a database connection", + // Repo constructors that open database connections; callers across a + // package boundary cannot be found by local taint propagation. + "github.com/smartcontractkit/chainlink-ccv/indexer/pkg/storage.NewPostgresStorage": "opens a database connection", + "github.com/smartcontractkit/chainlink-ccv/indexer/pkg/replay.NewStoreFromConfig": "opens a database connection", +} + +// denyMethods are methods that perform network/DB I/O, matched by name. An +// entry of the form "Name@substr" additionally requires the receiver type's +// string to contain substr (for generic names like Ping or Get). +var denyMethods = map[string]string{ + "Dial": "dials a remote endpoint", + "DialContext": "dials a remote endpoint", + "CallContext": "issues an RPC call", + "CallContract": "issues an RPC call", + "HeaderByNumber": "fetches a block header over RPC", + "HeaderByHash": "fetches a block header over RPC", + "BlockByNumber": "fetches a block over RPC", + "BlockByHash": "fetches a block over RPC", + "TransactionReceipt": "fetches a receipt over RPC", + "TransactionByHash": "fetches a transaction over RPC", + "SendTransaction": "sends a transaction over RPC", + "BalanceAt": "queries chain state over RPC", + "CodeAt": "queries chain state over RPC", + "NonceAt": "queries chain state over RPC", + "PendingNonceAt": "queries chain state over RPC", + "ChainID": "queries the chain over RPC", + "SyncProgress": "queries the chain over RPC", + "SuggestGasPrice": "queries the chain over RPC", + "SuggestGasTipCap": "queries the chain over RPC", + "EstimateGas": "queries the chain over RPC", + "FilterLogs": "queries logs over RPC", + "SubscribeFilterLogs": "subscribes to logs over RPC", + "LatestBlock": "fetches the latest block over RPC", + "LatestAndFinalizedBlock": "fetches blocks over RPC", + "LatestFinalizedBlock": "fetches the finalized block over RPC", + "LatestSafeBlock": "fetches the safe block over RPC", + "FetchMessageSentEvents": "fetches logs over RPC", + "GetRMNCursedSubjects": "reads an RMN Remote over RPC", + "GetStaticConfig": "reads contract state over RPC", + "GetDestChainConfig": "reads contract state over RPC", + "ReadChainStatuses": "reads from the database", + "WriteChainStatuses": "writes to the database", + "GetDiscoverySequenceNumber": "reads from the database", + "CreateDiscoveryState": "writes to the database", + "LoadJob": "reads from the database", + "RunMigrations": "runs database migrations", + "RunPostgresMigrations": "runs database migrations", + "EnsureDBConnection": "pings the database", + "EnabledAddresses": "reads from the keystore", + "GetKeys": "reads from the keystore", + "EnsureKey": "writes to the keystore", + "EnsureImportedKey": "writes to the keystore", + "LoadKeystore": "reads from the keystore", + "GetPublicKey": "queries the KMS", + "LoadKMSKeystore": "queries the KMS", + // Known repo entrypoints whose I/O happens across a package boundary and so + // cannot be found by local taint propagation. + "New@sqlutil/pg": "opens a database connection", + "CreateStorage@pkg/storage": "opens and migrates database storage", + "Ping@sql": "pings the database", + "PingContext@sql": "pings the database", + "Ping@redis": "pings Redis", + "PingContext@redis": "pings Redis", + "Do@net/http": "issues an HTTP request", + "Get@net/http": "issues an HTTP request", + "Post@net/http": "issues an HTTP request", + "Query@sql": "queries the database", + "QueryContext@sql": "queries the database", + "QueryRow@sql": "queries the database", + "QueryRowContext@sql": "queries the database", + "Exec@sql": "writes to the database", + "ExecContext@sql": "writes to the database", + "Select@sqlx": "queries the database", + "SelectContext@sqlx": "queries the database", + "Get@sqlx": "queries the database", + "GetContext@sqlx": "queries the database", +} + +// spawnMethods are methods that run a function literal on a background +// goroutine (e.g. sync.WaitGroup.Go, errgroup.Group.Go). I/O inside those +// literals does not block startup, so it is exempt. +var spawnMethods = map[string]bool{ + "Go": true, + "GoCtx": true, + "Spawn": true, + "GoN": true, + "GoSafe": true, +} + +// ioCall records a denylisted call site and why it was flagged. +type ioCall struct { + pos token.Pos + reason string +} + +// funcInfo accumulates what we know about one declared function or method. +type funcInfo struct { + decl *ast.FuncDecl + obj *types.Func + ioCalls []ioCall // direct denylisted calls (outside spawned goroutines) + callees map[*types.Func]bool // local functions called outside spawned goroutines + isStartup bool + startupKind string // "constructor" or "Start method" +} + +func run(pass *analysis.Pass) (any, error) { + insp := pass.ResultOf[inspect.Analyzer].(*inspector.Inspector) //nolint:errcheck,revive // guaranteed by RunDespiteErrors + Requires + if insp == nil { + return nil, nil + } + + infos := make(map[*types.Func]*funcInfo) + + // Pass 1: for every declared function, collect its direct denylisted calls + // and its edges to other in-package functions, skipping bodies that run on + // spawned goroutines. + nodeFilter := []ast.Node{(*ast.FuncDecl)(nil)} + insp.Preorder(nodeFilter, func(n ast.Node) { + decl := n.(*ast.FuncDecl) //nolint:errcheck,revive // nodeFilter only matches *ast.FuncDecl + if isTestFile(pass, decl.Pos()) { + return + } + obj, _ := pass.TypesInfo.Defs[decl.Name].(*types.Func) //nolint:revive // nil when the name declares no func + if obj == nil { + return + } + fi := &funcInfo{decl: decl, obj: obj, callees: make(map[*types.Func]bool)} + fi.isStartup, fi.startupKind = classifyStartup(decl) + walkBody(pass, decl.Body, fi) + infos[obj] = fi + }) + + // Pass 2: propagate I/O reachability through the local call graph until a + // fixed point. A function is tainted when it performs I/O itself or calls + // (synchronously) a tainted in-package function. + tainted := make(map[*types.Func]bool) + for obj, fi := range infos { + if len(fi.ioCalls) > 0 { + tainted[obj] = true + } + } + for changed := true; changed; { + changed = false + for obj, fi := range infos { + if tainted[obj] { + continue + } + for callee := range fi.callees { + if tainted[callee] { + tainted[obj] = true + changed = true + break + } + } + } + } + + // Pass 3: report. Startup functions get their direct I/O calls flagged, and + // calls into tainted local helpers flagged at the call site. + for _, fi := range infos { + if !fi.isStartup || !tainted[fi.obj] { + continue + } + for _, c := range fi.ioCalls { + pass.Reportf(c.pos, + "I/O in %s %s: %s; constructors and Start must not perform I/O — defer to query time (common/lazy) or a background goroutine that reports via Ready/HealthReport", + fi.startupKind, fi.decl.Name.Name, c.reason) + } + // Report calls into tainted helpers with their positions. + reportTaintedCallees(pass, fi, tainted, infos) + } + return nil, nil +} + +// reportTaintedCallees flags direct calls from a startup function to in-package +// helpers that (transitively) perform I/O. +func reportTaintedCallees(pass *analysis.Pass, fi *funcInfo, tainted map[*types.Func]bool, infos map[*types.Func]*funcInfo) { + if fi.decl.Body == nil { + return + } + ast.Inspect(fi.decl.Body, func(n ast.Node) bool { + switch node := n.(type) { + case *ast.GoStmt: + return false + case *ast.CallExpr: + if skipsCallBody(pass, node) { + return false + } + callee := localFunc(pass, node) + if callee != nil && callee != fi.obj && tainted[callee] { + reason := "performs I/O" + if len(infos[callee].ioCalls) > 0 { + reason = infos[callee].ioCalls[0].reason + } + if !suppressed(pass, fi.decl, node.Pos()) { + pass.Reportf(node.Pos(), + "%s %s calls %s, which performs I/O (%s); constructors and Start must not perform I/O — defer to query time (common/lazy) or a background goroutine that reports via Ready/HealthReport", + fi.startupKind, fi.decl.Name.Name, callee.Name(), reason) + } + } + } + return true + }) +} + +// classifyStartup reports whether decl is a startup path: an exported +// constructor (New*) or a method named Start. +func classifyStartup(decl *ast.FuncDecl) (bool, string) { + if decl.Recv != nil && decl.Name.Name == "Start" { + return true, "Start method" + } + if decl.Recv == nil && strings.HasPrefix(decl.Name.Name, "New") && decl.Name.IsExported() { + return true, "constructor" + } + return false, "" +} + +// walkBody inspects a function body for denylisted calls and local call edges. +// Bodies running on spawned goroutines (go statements, wg.Go(func(){...}), ...) +// are skipped: that I/O does not block startup. +func walkBody(pass *analysis.Pass, body *ast.BlockStmt, fi *funcInfo) { + if body == nil { + return + } + ast.Inspect(body, func(n ast.Node) bool { + switch node := n.(type) { + case *ast.GoStmt: + return false + case *ast.CallExpr: + if skipsCallBody(pass, node) { + return false + } + if reason, ok := denyReason(pass, node); ok { + // A suppressed call is dropped entirely: the justification + // covers it, so it must not taint callers up the chain. + if !suppressed(pass, fi.decl, node.Pos()) { + fi.ioCalls = append(fi.ioCalls, ioCall{pos: node.Pos(), reason: reason}) + } + } + if callee := localFunc(pass, node); callee != nil { + fi.callees[callee] = true + } + } + return true + }) +} + +// denyReason returns why a call is denylisted, if it is. +func denyReason(pass *analysis.Pass, call *ast.CallExpr) (string, bool) { + sel, ok := call.Fun.(*ast.SelectorExpr) + if !ok { + return "", false + } + fn, ok := pass.TypesInfo.ObjectOf(sel.Sel).(*types.Func) + if !ok { + return "", false + } + sig, _ := fn.Type().(*types.Signature) //nolint:revive // *types.Func always carries a *types.Signature + if sig == nil || sig.Recv() == nil { + // Package-level function. + if fn.Pkg() == nil { + return "", false + } + reason, ok := denyFuncs[fn.Pkg().Path()+"."+fn.Name()] + return reason, ok + } + // Method: match on name, plus an optional receiver-type substring. + for entry, reason := range denyMethods { + name, recvSubstr, _ := strings.Cut(entry, "@") + if fn.Name() != name { + continue + } + if recvSubstr == "" || strings.Contains(sig.Recv().Type().String(), recvSubstr) { + return reason, true + } + } + return "", false +} + +// localFunc resolves a call to a function or method declared in the current +// package, so taint can propagate across helper boundaries. +func localFunc(pass *analysis.Pass, call *ast.CallExpr) *types.Func { + var id *ast.Ident + switch fun := call.Fun.(type) { + case *ast.Ident: + id = fun + case *ast.SelectorExpr: + id = fun.Sel + default: + return nil + } + fn, ok := pass.TypesInfo.ObjectOf(id).(*types.Func) + if !ok || fn.Pkg() == nil || fn.Pkg() != pass.Pkg { + return nil + } + return fn +} + +// isSpawnCall reports whether the call runs a function literal on a background +// goroutine (e.g. wg.Go(func(){...}), errgroup.Go). +func isSpawnCall(pass *analysis.Pass, call *ast.CallExpr) bool { + sel, ok := call.Fun.(*ast.SelectorExpr) + if !ok || !spawnMethods[sel.Sel.Name] { + return false + } + for _, arg := range call.Args { + if _, isLit := arg.(*ast.FuncLit); isLit { + return true + } + } + return false +} + +// isDeferredRegistration reports whether the call merely registers a function +// literal for execution at query time (common/lazy.New). The closure body is +// not startup code: it runs on first use, and failures are not cached. +func isDeferredRegistration(pass *analysis.Pass, call *ast.CallExpr) bool { + sel, ok := call.Fun.(*ast.SelectorExpr) + if !ok { + return false + } + fn, ok := pass.TypesInfo.ObjectOf(sel.Sel).(*types.Func) + if !ok || fn.Pkg() == nil { + return false + } + if fn.Name() == "New" && fn.Pkg().Name() == "lazy" { + return true + } + return false +} + +// skipsCallBody reports whether the call's function-literal arguments must not +// be walked as part of the enclosing startup function: the closure either runs +// on a background goroutine or is deferred to query time. +func skipsCallBody(pass *analysis.Pass, call *ast.CallExpr) bool { + return isSpawnCall(pass, call) || isDeferredRegistration(pass, call) +} + +// suppressed reports whether a finding at pos carries a //nolint:noeagerio +// comment on its line, the line above, or the enclosing function's doc. +// Raw comment text is scanned because CommentGroup.Text() strips directives. +func suppressed(pass *analysis.Pass, decl *ast.FuncDecl, pos token.Pos) bool { + posn := pass.Fset.Position(pos) + if decl.Doc != nil && groupHasNolint(decl.Doc) { + return true + } + for _, file := range pass.Files { + for _, cg := range file.Comments { + end := pass.Fset.Position(cg.End()) + if end.Filename != posn.Filename { + continue + } + if (end.Line == posn.Line || end.Line == posn.Line-1) && groupHasNolint(cg) { + return true + } + } + } + return false +} + +func groupHasNolint(cg *ast.CommentGroup) bool { + for _, c := range cg.List { + if strings.Contains(c.Text, "nolint:noeagerio") { + return true + } + } + return false +} + +func isTestFile(pass *analysis.Pass, pos token.Pos) bool { + return strings.HasSuffix(pass.Fset.Position(pos).Filename, "_test.go") +} diff --git a/tools/noeagerio/analyzer_test.go b/tools/noeagerio/analyzer_test.go new file mode 100644 index 000000000..a4191f6ab --- /dev/null +++ b/tools/noeagerio/analyzer_test.go @@ -0,0 +1,11 @@ +package noeagerio + +import ( + "testing" + + "golang.org/x/tools/go/analysis/analysistest" +) + +func TestAnalyzer(t *testing.T) { + analysistest.Run(t, analysistest.TestData(), Analyzer, "a") +} diff --git a/tools/noeagerio/cmd/noeagerio/main.go b/tools/noeagerio/cmd/noeagerio/main.go new file mode 100644 index 000000000..047f74b10 --- /dev/null +++ b/tools/noeagerio/cmd/noeagerio/main.go @@ -0,0 +1,17 @@ +// Command noeagerio is a vettool that reports I/O in constructors and Start +// methods. Run it via go vet: +// +// go vet -vettool=$(which noeagerio) ./... +// +// See the noeagerio package docs for the rule and its rationale. +package main + +import ( + "golang.org/x/tools/go/analysis/singlechecker" + + "github.com/smartcontractkit/chainlink-ccv/tools/noeagerio" +) + +func main() { + singlechecker.Main(noeagerio.Analyzer) +} diff --git a/tools/noeagerio/testdata/src/a/a.go b/tools/noeagerio/testdata/src/a/a.go new file mode 100644 index 000000000..bf8047cd9 --- /dev/null +++ b/tools/noeagerio/testdata/src/a/a.go @@ -0,0 +1,99 @@ +package a + +import ( + "context" + "net/http" + "sync" + + "common/lazy" +) + +// dialer is a stand-in for an RPC client; Dial is on the denylist by name. +type dialer struct{} + +func (d *dialer) Dial(ctx context.Context) error { return nil } + +type service struct { + d *dialer + wg sync.WaitGroup +} + +// Constructors that only assemble dependencies are fine. +func NewServiceClean(d *dialer) *service { + return &service{d: d} +} + +func NewServiceHTTP() *service { + _, _ = http.Get("http://example.com") // want `I/O in constructor NewServiceHTTP: issues an HTTP request` + return &service{} +} + +func NewServiceDial() *service { + s := &service{d: &dialer{}} + _ = s.d.Dial(context.Background()) // want `I/O in constructor NewServiceDial: dials a remote endpoint` + return s +} + +// helper does I/O but is not itself a startup path, so it is not flagged +// directly; callers in startup paths are flagged instead. +func (s *service) helper() error { + return s.d.Dial(context.Background()) +} + +// Start that calls a tainted helper is flagged at the call site. +func (s *service) Start(ctx context.Context) error { + return s.helper() // want `Start method Start calls .*helper, which performs I/O` +} + +// cleanService shows the blessed pattern: I/O happens in a spawned goroutine. +type cleanService struct { + d *dialer + wg sync.WaitGroup +} + +func (s *cleanService) Start(ctx context.Context) error { + s.wg.Go(func() { + _ = s.d.Dial(ctx) + }) + go func() { + _ = s.d.Dial(ctx) + }() + return nil +} + +// nolintService documents a deliberate exception. +type nolintService struct { + d *dialer +} + +func (s *nolintService) Start(ctx context.Context) error { + //nolint:noeagerio // fail-fast by design: this dependency is required for anything to work + return s.d.Dial(ctx) +} + +// unexported constructors are not flagged (callers own the policy), and +// neither are non-startup functions. +func newUnexported(d *dialer) *service { + _ = d.Dial(context.Background()) + return &service{d: d} +} + +func queryTime(ctx context.Context, d *dialer) error { + _, err := http.Get("http://example.com") + _ = d.Dial(ctx) + return err +} + +// lazyService defers its I/O to query time via common/lazy — the blessed +// pattern for values that need a derivation RPC. +type lazyService struct { + addr *lazy.Lazy[string] +} + +func NewLazyService(d *dialer) *lazyService { + return &lazyService{ + addr: lazy.New(func(ctx context.Context) (string, error) { + return "", d.Dial(ctx) + }), + } +} diff --git a/tools/noeagerio/testdata/src/common/lazy/lazy.go b/tools/noeagerio/testdata/src/common/lazy/lazy.go new file mode 100644 index 000000000..9c8745175 --- /dev/null +++ b/tools/noeagerio/testdata/src/common/lazy/lazy.go @@ -0,0 +1,17 @@ +// Package lazy mirrors github.com/smartcontractkit/chainlink-ccv/common/lazy +// for analyzer tests: New only registers the derivation for query time. +package lazy + +import "context" + +type Lazy[V any] struct { + derive func(ctx context.Context) (V, error) +} + +func New[V any](derive func(ctx context.Context) (V, error)) *Lazy[V] { + return &Lazy[V]{derive: derive} +} + +func (l *Lazy[V]) Value(ctx context.Context) (V, error) { + return l.derive(ctx) +} diff --git a/verifier/pkg/commit/signer.go b/verifier/pkg/commit/signer.go index c4b4986bc..f63a6d188 100644 --- a/verifier/pkg/commit/signer.go +++ b/verifier/pkg/commit/signer.go @@ -205,6 +205,7 @@ func (a *KeystoreSignerAdapter) Sign(data []byte) ([]byte, error) { // along with the public key (as UnknownAddress) for use in committee configuration. func NewSignerFromKeystore(ctx context.Context, ks keystore.Keystore, keyName string) (signer verifier.MessageSigner, pubKey []byte, address protocol.UnknownAddress, err error) { // Get the key info to retrieve the public key + //nolint:noeagerio // fail-fast by design: a verifier without its signing key can do no work, so the keystore read happens at startup rather than at first sign keysResp, err := ks.GetKeys(ctx, keystore.GetKeysRequest{KeyNames: []string{keyName}}) if err != nil { return nil, nil, protocol.UnknownAddress{}, fmt.Errorf("failed to get key from keystore: %w", err) diff --git a/verifier/pkg/coordinator.go b/verifier/pkg/coordinator.go index 44d7c6d00..140ca8a82 100644 --- a/verifier/pkg/coordinator.go +++ b/verifier/pkg/coordinator.go @@ -51,9 +51,9 @@ type Coordinator struct { initFn func(ctx context.Context) error - curseDetector common.CurseCheckerService - chainStatusBatcher *chainstatus.Batcher - sourceReaderServices map[protocol.ChainSelector]services.Service + curseDetector common.CurseCheckerService + chainStatusBatcher *chainstatus.Batcher + sourceReaderServices map[protocol.ChainSelector]services.Service // sourceReaderStartErrs records per-chain source readers that failed to // start so HealthReport keeps them visible after they are skipped. sourceReaderStartErrs map[protocol.ChainSelector]error @@ -134,11 +134,11 @@ func NewCoordinatorWithDetector( } lggr = logger.With(lggr, "verifierID", config.VerifierID) vc := &Coordinator{ - lggr: lggr, - verifierID: config.VerifierID, - monitoring: monitoring, - messageRulesSvc: messageRulesSvc, - sourceReaderStartErrs: make(map[protocol.ChainSelector]error), + lggr: lggr, + verifierID: config.VerifierID, + monitoring: monitoring, + messageRulesSvc: messageRulesSvc, + sourceReaderStartErrs: make(map[protocol.ChainSelector]error), } vc.initFn = func(ctx context.Context) error { // Batch the chain status writes. The source readers write a status on every @@ -447,9 +447,14 @@ func filterConfiguredSourceReaders( allSelectors = append(allSelectors, selector) } - statusMap, err := chainStatusManager.ReadChainStatuses(ctx, allSelectors) + // A transient DB failure here must not abort startup: the statuses only + // inform logging and the disabled-chain warning below. Degrade to unknown + // statuses; each source reader re-reads its own status (with retries) when + // it initializes its start block. + statusMap, err := chainStatusManager.ReadChainStatuses(ctx, allSelectors) //nolint:noeagerio // single bounded read of the service's own DB at coordinator start; failure is non-fatal if err != nil { - return nil, fmt.Errorf("failed to read chain statuses from storage: %w", err) + lggr.Errorw("Failed to read chain statuses from storage, continuing with unknown statuses", "error", err) + statusMap = nil } configuredSourceReaders := make(map[protocol.ChainSelector]chainaccess.SourceReader) diff --git a/verifier/pkg/coordinator_start_test.go b/verifier/pkg/coordinator_start_test.go index 36bfb7c2a..1fa671203 100644 --- a/verifier/pkg/coordinator_start_test.go +++ b/verifier/pkg/coordinator_start_test.go @@ -31,10 +31,12 @@ func (f *fakeSourceReaderService) Start(context.Context) error { return nil } -func (f *fakeSourceReaderService) Close() error { return nil } -func (f *fakeSourceReaderService) Name() string { return f.name } -func (f *fakeSourceReaderService) Ready() error { return nil } -func (f *fakeSourceReaderService) HealthReport() map[string]error { return map[string]error{f.name: nil} } +func (f *fakeSourceReaderService) Close() error { return nil } +func (f *fakeSourceReaderService) Name() string { return f.name } +func (f *fakeSourceReaderService) Ready() error { return nil } +func (f *fakeSourceReaderService) HealthReport() map[string]error { + return map[string]error{f.name: nil} +} // One chain failing to start must not stop the others, and the failure must // stay visible in the health report. From 30608bda2dac00c3a8f1e5650838ae9cd0458dc0 Mon Sep 17 00:00:00 2001 From: Terry Tata Date: Wed, 7 Oct 2026 09:43:34 -0700 Subject: [PATCH 3/3] fix: bound startup DB I/O by caller ctx; build chain accessors concurrently Audit follow-ups: - EnsureDBConnectionContext / RunPostgresMigrationsContext / RunMigrationsContext: ping retries and goose migrations now honor the caller's context (previously ~40s of unbounded retry, and migrations with no ctx at all). Old signatures kept as deprecated wrappers. - ConnectToPostgresDB takes a ctx so a degraded Postgres cannot blow the verifier startup budget. - Verifier (committee + token) and executor factories build chain accessors concurrently with a 30s per-chain timeout: a slow RPC pool no longer serializes away the shared startup budget. - messagerules initial poll uses the service-lifetime ctx instead of the startup ctx bootstrap cancels on return. - Coordinator startup chain-status read is bounded (5s) and non-fatal. - noeagerio denylist gains GetAccessor (accessor construction dials). --- aggregator/cmd/main.go | 9 +-- .../read_path_filter_config_change_test.go | 2 +- aggregator/pkg/server.go | 9 +-- aggregator/pkg/storage/factory.go | 12 ++-- .../pkg/storage/postgres/run_migrations.go | 12 +++- aggregator/tests/utils.go | 4 +- bootstrap/bootstrap.go | 2 +- bootstrap/bootstrap_test.go | 2 +- bootstrap/db/run_migrations.go | 13 +++- changelog/2026-10-06_no_eager_io_startup.md | 24 ++++++- cmd/executor/service.go | 60 ++++++++++++------ cmd/verifier/common.go | 10 ++- cmd/verifier/run_ccv_cli.go | 7 ++- cmd/verifier/servicefactory.go | 63 ++++++++++++------- cmd/verifier/tokenfactory.go | 47 +++++++++----- common/db_utils.go | 31 ++++++--- .../jobqueue/postgres_queue_explain_test.go | 2 +- indexer/cmd/main.go | 4 +- indexer/cmd/replay/main.go | 4 +- indexer/pkg/storage/run_migrations.go | 12 +++- integration/pkg/messagerules/poller.go | 7 ++- tools/noeagerio/analyzer.go | 1 + tools/noeagerio/testdata/src/a/a.go | 17 +++++ verifier/pkg/coordinator.go | 12 ++-- verifier/pkg/db/run_migrations.go | 12 +++- verifier/testutil/test_db.go | 2 +- 26 files changed, 275 insertions(+), 105 deletions(-) diff --git a/aggregator/cmd/main.go b/aggregator/cmd/main.go index 7fd9d0b65..6499aa3cd 100644 --- a/aggregator/cmd/main.go +++ b/aggregator/cmd/main.go @@ -210,14 +210,15 @@ func runServer(configPath, logLevelStr string, lggr logger.Logger, sugaredLggr l protocol.InitChainSelectorCache() - server, err := aggregator.NewServer(sugaredLggr, config, aggMonitoring) - if err != nil { - sugaredLggr.Fatalw("failed to create CCV data service", "error", err) - } ctx := context.Background() ctx, stop := signal.NotifyContext(ctx, os.Interrupt, syscall.SIGTERM) defer stop() + server, err := aggregator.NewServer(ctx, sugaredLggr, config, aggMonitoring) + if err != nil { + sugaredLggr.Fatalw("failed to create CCV data service", "error", err) + } + lc := &net.ListenConfig{} lis, err := lc.Listen(ctx, "tcp", config.Server.Address) if err != nil { diff --git a/aggregator/pkg/aggregation/read_path_filter_config_change_test.go b/aggregator/pkg/aggregation/read_path_filter_config_change_test.go index 7660d4f05..db30faee2 100644 --- a/aggregator/pkg/aggregation/read_path_filter_config_change_test.go +++ b/aggregator/pkg/aggregation/read_path_filter_config_change_test.go @@ -20,7 +20,7 @@ import ( func setupTestPostgresStorage(t *testing.T) (*postgres.DatabaseStorage, func()) { t.Helper() ds, cleanup := testutil.SetupTestPostgresDB(t) - err := postgres.RunMigrations(ds, "postgres") + err := postgres.RunMigrationsContext(t.Context(), ds, "postgres") require.NoError(t, err) storage := postgres.NewDatabaseStorage(ds, 10, 10*time.Second, logger.Sugared(logger.Test(t))) return storage, cleanup diff --git a/aggregator/pkg/server.go b/aggregator/pkg/server.go index f9992f8ce..7fa922be8 100644 --- a/aggregator/pkg/server.go +++ b/aggregator/pkg/server.go @@ -310,8 +310,9 @@ type SignatureAndQuorumValidator interface { // NewServer creates a new aggregator server with the specified logger, configuration, and monitoring. // aggMonitoring must not be nil; use monitoring.NoopAggregatorMonitoring when monitoring is disabled. // Errors are returned to the caller (main), which owns the fail-fast decision; a -// library constructor must never exit the process. -func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonitoring common.AggregatorMonitoring) (*Server, error) { +// library constructor must never exit the process. ctx bounds the startup DB +// connect ping and migrations. +func NewServer(ctx context.Context, l logger.SugaredLogger, config *model.AggregatorConfig, aggMonitoring common.AggregatorMonitoring) (*Server, error) { if err := config.Validate(); err != nil { return nil, fmt.Errorf("failed to validate server configuration: %w", err) } @@ -327,8 +328,8 @@ func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonito ) factory := storage.NewStorageFactory(l) - //nolint:noeagerio // the aggregator's own DB is a hard dependency: connect + migrate fail fast at startup, and the health endpoint reports readiness after that - rawStore, err := factory.CreateStorage(config.Storage, aggMonitoring) + //nolint:noeagerio // the aggregator's own DB is a hard dependency: connect + migrate fail fast at startup (bounded by ctx), and the health endpoint reports readiness after that + rawStore, err := factory.CreateStorage(ctx, config.Storage, aggMonitoring) if err != nil { return nil, fmt.Errorf("failed to create storage: %w", err) } diff --git a/aggregator/pkg/storage/factory.go b/aggregator/pkg/storage/factory.go index cea5b192a..03288622d 100644 --- a/aggregator/pkg/storage/factory.go +++ b/aggregator/pkg/storage/factory.go @@ -1,6 +1,7 @@ package storage import ( + "context" "database/sql" "fmt" "time" @@ -48,15 +49,16 @@ func NewStorageFactory(logger logger.SugaredLogger) *Factory { } // CreateStorage creates a storage instance based on the provided configuration. -func (f *Factory) CreateStorage(config *model.StorageConfig, monitoring common.AggregatorMonitoring) (CommitVerificationStorage, error) { +// The connect ping and migrations are bounded by ctx. +func (f *Factory) CreateStorage(ctx context.Context, config *model.StorageConfig, monitoring common.AggregatorMonitoring) (CommitVerificationStorage, error) { if config.StorageType != model.StorageTypePostgreSQL { return nil, fmt.Errorf("unsupported storage type: %s (only postgres is supported)", config.StorageType) } - return f.createPostgreSQLStorage(config) + return f.createPostgreSQLStorage(ctx, config) } // createPostgreSQLStorage creates a PostgreSQL-backed storage instance. -func (f *Factory) createPostgreSQLStorage(config *model.StorageConfig) (CommitVerificationStorage, error) { +func (f *Factory) createPostgreSQLStorage(ctx context.Context, config *model.StorageConfig) (CommitVerificationStorage, error) { if config.ConnectionURL == "" { return nil, fmt.Errorf("PostgreSQL connection URL is required") } @@ -98,7 +100,7 @@ func (f *Factory) createPostgreSQLStorage(config *model.StorageConfig) (CommitVe "connMaxIdleTime", connMaxIdleTime, ) - if err := ccvcommon.EnsureDBConnection(f.logger, db); err != nil { + if err := ccvcommon.EnsureDBConnectionContext(ctx, f.logger, db); err != nil { return nil, fmt.Errorf("failed to ping PostgreSQL database: %w", err) } @@ -106,7 +108,7 @@ func (f *Factory) createPostgreSQLStorage(config *model.StorageConfig) (CommitVe sqlxDB := sqlx.NewDb(db, postgresDriver) // Run PostgreSQL migrations - err = postgres.RunMigrations(sqlxDB, postgresDriver) + err = postgres.RunMigrationsContext(ctx, sqlxDB, postgresDriver) if err != nil { return nil, fmt.Errorf("failed to run PostgreSQL migrations: %w", err) } diff --git a/aggregator/pkg/storage/postgres/run_migrations.go b/aggregator/pkg/storage/postgres/run_migrations.go index d12a531a6..0ee57da5a 100644 --- a/aggregator/pkg/storage/postgres/run_migrations.go +++ b/aggregator/pkg/storage/postgres/run_migrations.go @@ -1,6 +1,7 @@ package postgres import ( + "context" "fmt" "sync" @@ -13,7 +14,16 @@ import ( var migrationMutex = sync.Mutex{} // RunMigrations applies database-specific SQL migrations. +// +// Deprecated: use RunMigrationsContext so a caller's deadline can abort a hung +// migration; this wrapper is unbounded by any caller context. func RunMigrations(db *sqlx.DB, dbType string) error { + return RunMigrationsContext(context.Background(), db, dbType) +} + +// RunMigrationsContext applies PostgreSQL database migrations, aborting when +// ctx is done. +func RunMigrationsContext(ctx context.Context, db *sqlx.DB, dbType string) error { migrationMutex.Lock() defer migrationMutex.Unlock() @@ -30,7 +40,7 @@ func RunMigrations(db *sqlx.DB, dbType string) error { return fmt.Errorf("failed to set goose dialect: %w", err) } - if err := goose.Up(db.DB, "postgres"); err != nil { + if err := goose.UpContext(ctx, db.DB, "postgres"); err != nil { return fmt.Errorf("failed to run postgres migrations: %w", err) } diff --git a/aggregator/tests/utils.go b/aggregator/tests/utils.go index f1ed8f8c7..821482ca3 100644 --- a/aggregator/tests/utils.go +++ b/aggregator/tests/utils.go @@ -222,7 +222,7 @@ func CreateServerOnlyWithMessageRulesControl(t *testing.T, options ...ConfigOpti } config.Storage = storageConfig - rawRuleStore, err := storage.NewStorageFactory(sugaredLggr).CreateStorage(config.Storage, monitoring.NewNoopAggregatorMonitoring()) + rawRuleStore, err := storage.NewStorageFactory(sugaredLggr).CreateStorage(t.Context(), config.Storage, monitoring.NewNoopAggregatorMonitoring()) if err != nil { cleanupStorage() return nil, nil, nil, err @@ -233,7 +233,7 @@ func CreateServerOnlyWithMessageRulesControl(t *testing.T, options ...ConfigOpti return nil, nil, nil, fmt.Errorf("test storage does not implement message rules store") } - s, err := agg.NewServer(sugaredLggr, config, monitoring.NewNoopAggregatorMonitoring()) + s, err := agg.NewServer(t.Context(), sugaredLggr, config, monitoring.NewNoopAggregatorMonitoring()) if err != nil { cleanupStorage() return nil, nil, nil, fmt.Errorf("failed to create server: %w", err) diff --git a/bootstrap/bootstrap.go b/bootstrap/bootstrap.go index 45d1275bb..714d03980 100644 --- a/bootstrap/bootstrap.go +++ b/bootstrap/bootstrap.go @@ -774,7 +774,7 @@ func connectToDB(ctx context.Context, connStr string) (*sqlx.DB, error) { if err != nil { return nil, fmt.Errorf("failed to connect to bootstrapper database: %w", err) } - if err := dbpkg.RunMigrations(db); err != nil { + if err := dbpkg.RunMigrationsContext(ctx, db); err != nil { return nil, fmt.Errorf("failed to run bootstrapper database migrations: %w", err) } return db, nil diff --git a/bootstrap/bootstrap_test.go b/bootstrap/bootstrap_test.go index c809c1366..ae990c833 100644 --- a/bootstrap/bootstrap_test.go +++ b/bootstrap/bootstrap_test.go @@ -102,7 +102,7 @@ func TestBootstrapDB_RunMigrations(t *testing.T) { require.NoError(t, err) defer dbConn.Close() - err = db.RunMigrations(dbConn) + err = db.RunMigrationsContext(ctx, dbConn) require.NoError(t, err) var count int diff --git a/bootstrap/db/run_migrations.go b/bootstrap/db/run_migrations.go index 3b36e2bf9..23927588c 100644 --- a/bootstrap/db/run_migrations.go +++ b/bootstrap/db/run_migrations.go @@ -1,6 +1,7 @@ package db import ( + "context" "embed" "fmt" "sync" @@ -14,7 +15,17 @@ var migrations embed.FS var migrationMutex = sync.Mutex{} +// RunMigrations applies the bootstrap database migrations. +// +// Deprecated: use RunMigrationsContext so a caller's startup deadline can abort +// a hung migration; this wrapper is unbounded by any caller context. func RunMigrations(db *sqlx.DB) error { + return RunMigrationsContext(context.Background(), db) +} + +// RunMigrationsContext applies the bootstrap database migrations, aborting when +// ctx is done. +func RunMigrationsContext(ctx context.Context, db *sqlx.DB) error { migrationMutex.Lock() defer migrationMutex.Unlock() @@ -24,7 +35,7 @@ func RunMigrations(db *sqlx.DB) error { return fmt.Errorf("failed to set goose dialect: %w", err) } - if err := goose.Up(db.DB, "migrations"); err != nil { + if err := goose.UpContext(ctx, db.DB, "migrations"); err != nil { return fmt.Errorf("failed to run migrations: %w", err) } diff --git a/changelog/2026-10-06_no_eager_io_startup.md b/changelog/2026-10-06_no_eager_io_startup.md index efb9804a5..522d58940 100644 --- a/changelog/2026-10-06_no_eager_io_startup.md +++ b/changelog/2026-10-06_no_eager_io_startup.md @@ -20,14 +20,34 @@ repeated verifier outages (see also PR #1424): * `pricer`: a chain that fails to start is skipped and surfaced via the new `HealthReport()`; fails only when no chain started. * `aggregator/pkg`: `NewServer` returns errors instead of calling - `logger.Fatalf` (signature changed to `(*Server, error)`); the caller in - `main` owns the fail-fast decision. + `logger.Fatalf` (signature changed to `(ctx, ...) (*Server, error)`); the + caller in `main` owns the fail-fast decision. +* Startup DB work is now bounded by the caller's context: + `common.EnsureDBConnectionContext` (the old `EnsureDBConnection` retried for + ~40s ignoring the startup deadline) and `RunPostgresMigrationsContext` / + `RunMigrationsContext` (goose `UpContext`; a hung migration no longer hangs + startup forever). `cmd/verifier.ConnectToPostgresDB` takes a ctx. +* Accessor construction (RPC dial + chain service/TXM start) runs concurrently + per chain with a 30s per-chain timeout in the committee verifier, token + verifier, and executor factories, so one slow chain can no longer serialize + away the shared startup budget; failures remain skip-and-log. +* `integration/pkg/messagerules`: the initial rules poll uses the + service-lifetime context instead of the startup ctx that bootstrap cancels + as soon as `Start` returns (the first fetch could previously be aborted). +* The coordinator's startup chain-status read is bounded (5s) in addition to + being non-fatal. * New `noeagerio` static analyzer (`tools/noeagerio`, run by `just lint-noeagerio` and the lint CI workflow) forbids I/O in constructors and `Start` methods repo-wide, with `//nolint:noeagerio` as the documented escape hatch for deliberate fail-fast exceptions. Policy recorded in AGENTS.md. +Verified non-blocking by inspection: pyroscope `Start` (async uploader +goroutines) and beholder `SetupBeholder` (lazy gRPC client). Still structural +and accepted: JD-mode cached-job startup runs inside the lifecycle manager's +`Start` (its worst case shrinks with every fix above); a larger redesign would +move the startup deadline from the factory call to the readiness probe. + Deliberately unchanged (fail-fast by design, annotated with `//nolint:noeagerio` + justification): bootstrap DB connect/migrations and keystore/KMS key verification, the JD lifecycle cached-job load, the commit signer keystore diff --git a/cmd/executor/service.go b/cmd/executor/service.go index 065b94ba6..5a4d082ec 100644 --- a/cmd/executor/service.go +++ b/cmd/executor/service.go @@ -6,11 +6,13 @@ import ( "fmt" "net/http" "strconv" + "sync" "time" "github.com/gin-gonic/gin" "github.com/grafana/pyroscope-go" sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "golang.org/x/sync/errgroup" "github.com/smartcontractkit/chainlink-ccv/bootstrap" "github.com/smartcontractkit/chainlink-ccv/common/health" @@ -36,6 +38,10 @@ const ( indexerGarbageCollectionInterval = 1 * time.Hour // httpShutdownTimeout bounds the graceful drain of the HTTP server during Stop. httpShutdownTimeout = 5 * time.Second + // accessorBuildTimeout caps one chain's accessor construction (RPC dial + + // TXM start) during factory Start; the parent startup context may be + // tighter and always wins. + accessorBuildTimeout = 30 * time.Second ) // Factory is a bootstrap.ServiceFactory that starts the executor service. @@ -142,6 +148,11 @@ func (f *Factory) Start(ctx context.Context, spec bootstrap.JobSpec, deps bootst rmnReaders := make(map[protocol.ChainSelector]chainaccess.RMNCurseReader) enabledDestChains := make([]protocol.ChainSelector, 0) + // Chains are built concurrently with a per-chain timeout: GetAccessor dials + // the chain and starts its TXM, so a slow (not failing) chain must not + // serialize away the shared startup budget. + var chainsMu sync.Mutex + g, gctx := errgroup.WithContext(ctx) for strSel := range executorConfig.ChainConfiguration { selectorUint, err := strconv.ParseUint(strSel, 10, 64) if err != nil { @@ -150,27 +161,36 @@ func (f *Factory) Start(ctx context.Context, spec bootstrap.JobSpec, deps bootst } selector := protocol.ChainSelector(selectorUint) - accessor, err := deps.Registry.GetAccessor(ctx, selector) - if err != nil { - f.lggr.Errorw("Failed to get accessor for chain", "error", err, "chainSelector", strSel) - continue - } - - dr, drErr := accessor.DestinationReader() - ct, ctErr := accessor.ContractTransmitter() - - if drErr != nil || ctErr != nil { - f.lggr.Warnw("Skipping chain: missing DestinationReader or ContractTransmitter", "chainSelector", strSel, "destReaderErr", drErr, "transmitterErr", ctErr) - continue - } - - attachExecutorMonitoring(dr, ct, executorMonitoring) - - destReaders[selector] = dr - rmnReaders[selector] = dr - contractTransmitters[selector] = ct - enabledDestChains = append(enabledDestChains, selector) + g.Go(func() error { + chainCtx, cancel := context.WithTimeout(gctx, accessorBuildTimeout) + defer cancel() + //nolint:noeagerio // accessor construction dials the chain and starts its TXM; it runs concurrently per chain with a bounded ctx, and failures are skipped, not fatal + accessor, err := deps.Registry.GetAccessor(chainCtx, selector) + if err != nil { + f.lggr.Errorw("Failed to get accessor for chain", "error", err, "chainSelector", strSel) + return nil + } + + dr, drErr := accessor.DestinationReader() + ct, ctErr := accessor.ContractTransmitter() + + if drErr != nil || ctErr != nil { + f.lggr.Warnw("Skipping chain: missing DestinationReader or ContractTransmitter", "chainSelector", strSel, "destReaderErr", drErr, "transmitterErr", ctErr) + return nil + } + + attachExecutorMonitoring(dr, ct, executorMonitoring) + + chainsMu.Lock() + defer chainsMu.Unlock() + destReaders[selector] = dr + rmnReaders[selector] = dr + contractTransmitters[selector] = ct + enabledDestChains = append(enabledDestChains, selector) + return nil + }) } + _ = g.Wait() curseChecker := cursechecker.NewCachedCurseChecker(cursechecker.Params{ Lggr: f.lggr, diff --git a/cmd/verifier/common.go b/cmd/verifier/common.go index a2de02635..cc0926d8a 100644 --- a/cmd/verifier/common.go +++ b/cmd/verifier/common.go @@ -1,6 +1,7 @@ package verifier import ( + "context" "database/sql" "fmt" "time" @@ -30,7 +31,10 @@ const ( // resolved from the verifier secrets file when present, otherwise from CL_DATABASE_URL (the file // wins). A nil secrets argument is the env-only path. An empty URL leaves the DB // unconfigured (returns a nil DataSource), preserving the existing "DB optional" behavior. -func ConnectToPostgresDB(lggr logger.Logger, secrets *vsecrets.VerifierSecrets) (sqlutil.DataSource, error) { +// +// The connect ping retries and the migrations are bounded by ctx, so a degraded +// database cannot block startup past the caller's deadline. +func ConnectToPostgresDB(ctx context.Context, lggr logger.Logger, secrets *vsecrets.VerifierSecrets) (sqlutil.DataSource, error) { dbURL := secrets.DatabaseURL() if dbURL == "" { return nil, nil @@ -46,14 +50,14 @@ func ConnectToPostgresDB(lggr logger.Logger, secrets *vsecrets.VerifierSecrets) dbx.SetConnMaxLifetime(defaultConnMaxLifetime) dbx.SetConnMaxIdleTime(defaultConnMaxIdleTime) - if err := ccvcommon.EnsureDBConnection(lggr, dbx); err != nil { + if err := ccvcommon.EnsureDBConnectionContext(ctx, lggr, dbx); err != nil { _ = dbx.Close() return nil, fmt.Errorf("failed to ping postgres database: %w", err) } sqlxDB := sqlx.NewDb(dbx, "postgres") - if err := db.RunPostgresMigrations(sqlxDB); err != nil { + if err := db.RunPostgresMigrationsContext(ctx, sqlxDB); err != nil { _ = dbx.Close() return nil, fmt.Errorf("failed to run postgres migrations: %w", err) } diff --git a/cmd/verifier/run_ccv_cli.go b/cmd/verifier/run_ccv_cli.go index 7f492597b..aaf8b3ba9 100644 --- a/cmd/verifier/run_ccv_cli.go +++ b/cmd/verifier/run_ccv_cli.go @@ -1,6 +1,7 @@ package verifier import ( + "context" "fmt" "os" "path/filepath" @@ -58,7 +59,7 @@ func RunCCVCLI(args []string, secretsEnvVar, defaultSecretsPath string) { var chainStatusDeps chainstatuses.Deps getChainStatusDeps := func() chainstatuses.Deps { chainStatusOnce.Do(func() { - ds, connErr := ConnectToPostgresDB(lggr, secrets) + ds, connErr := ConnectToPostgresDB(context.Background(), lggr, secrets) if connErr != nil { _, _ = fmt.Fprintf(os.Stderr, "failed to connect to database: %v\n", connErr) os.Exit(1) @@ -77,7 +78,7 @@ func RunCCVCLI(args []string, secretsEnvVar, defaultSecretsPath string) { var jobQueueDeps jobqueue.Deps getJobQueueDeps := func() jobqueue.Deps { jobQueueOnce.Do(func() { - ds, connErr := ConnectToPostgresDB(lggr, secrets) + ds, connErr := ConnectToPostgresDB(context.Background(), lggr, secrets) if connErr != nil { _, _ = fmt.Fprintf(os.Stderr, "failed to connect to database: %v\n", connErr) os.Exit(1) @@ -96,7 +97,7 @@ func RunCCVCLI(args []string, secretsEnvVar, defaultSecretsPath string) { var recoveryStore recoverycli.Store getRecoveryStore := func() recoverycli.Store { recoveryOnce.Do(func() { - ds, err := ConnectToPostgresDB(lggr, secrets) + ds, err := ConnectToPostgresDB(context.Background(), lggr, secrets) if err != nil || ds == nil { _, _ = fmt.Fprintf(os.Stderr, "recovery requires a database connection: %v\n", err) os.Exit(1) diff --git a/cmd/verifier/servicefactory.go b/cmd/verifier/servicefactory.go index 5fa801feb..9f70c4fc5 100644 --- a/cmd/verifier/servicefactory.go +++ b/cmd/verifier/servicefactory.go @@ -7,11 +7,13 @@ import ( "net/http" "slices" "strconv" + "sync" "time" "github.com/gin-gonic/gin" "github.com/grafana/pyroscope-go" sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "golang.org/x/sync/errgroup" "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" @@ -41,7 +43,13 @@ import ( const defaultHTTPListenPort = 8100 // httpShutdownTimeout bounds the graceful drain of the HTTP server during Stop. -const httpShutdownTimeout = 5 * time.Second +const ( + httpShutdownTimeout = 5 * time.Second + // accessorBuildTimeout caps one chain's accessor construction (RPC dial + + // chain services) during factory Start; the parent startup context may be + // tighter and always wins. + accessorBuildTimeout = 30 * time.Second +) // factory is a ServiceFactory implementation that creates a committee verifier service. type factory struct { @@ -211,27 +219,40 @@ func (f *factory) Start(ctx context.Context, spec bootstrap.JobSpec, deps bootst // A failure to stand up one chain's reader (e.g. an unreachable RPC) must not stop the // remaining chains from starting. Log and skip, then only reject the whole coordinator if no - // chain is usable. + // chain is usable. Chains are built concurrently with a per-chain timeout: accessor + // construction dials the chain, so a slow (not failing) RPC must not serialize away the + // shared startup budget. chainSelectors := chainaccess.Infos[string](config.OnRampAddresses).GetAllChainSelectors() sourceReaders := make(map[protocol.ChainSelector]chainaccess.SourceReader) + var readersMu sync.Mutex + g, gctx := errgroup.WithContext(ctx) for _, selector := range chainSelectors { - accessor, err := deps.Registry.GetAccessor(ctx, selector) - if err != nil { - lggr.Errorw("Failed to get accessor, skipping chain", "error", err, "selector", selector) - continue - } - reader, err := accessor.SourceReader() - if err != nil { - lggr.Errorw("Failed to get source reader, skipping chain", "selector", selector, "error", err) - continue - } - observedReader, err := instrumentSourceReader(reader, config.VerifierID, selector, verifierMonitoring) - if err != nil { - lggr.Errorw("Failed to instrument source reader, skipping chain", "selector", selector, "error", err) - continue - } - sourceReaders[selector] = observedReader + g.Go(func() error { + chainCtx, cancel := context.WithTimeout(gctx, accessorBuildTimeout) + defer cancel() + //nolint:noeagerio // accessor construction dials the chain; it runs concurrently per chain with a bounded ctx, and failures are skipped, not fatal + accessor, err := deps.Registry.GetAccessor(chainCtx, selector) + if err != nil { + lggr.Errorw("Failed to get accessor, skipping chain", "error", err, "selector", selector) + return nil + } + reader, err := accessor.SourceReader() + if err != nil { + lggr.Errorw("Failed to get source reader, skipping chain", "selector", selector, "error", err) + return nil + } + observedReader, err := instrumentSourceReader(reader, config.VerifierID, selector, verifierMonitoring) + if err != nil { + lggr.Errorw("Failed to instrument source reader, skipping chain", "selector", selector, "error", err) + return nil + } + readersMu.Lock() + sourceReaders[selector] = observedReader + readersMu.Unlock() + return nil + }) } + _ = g.Wait() if len(sourceReaders) == 0 { return fmt.Errorf("no source readers configured: ensure at least one chain has a working source reader") } @@ -278,7 +299,7 @@ func (f *factory) Start(ctx context.Context, spec bootstrap.JobSpec, deps bootst lggr.Infow("Using signer address", "address", signerAddress) // Create chain status manager (PostgreSQL storage) with monitoring decorator - chainStatusManager, chainStatusDB, err := createChainStatusManager(lggr, config.VerifierID, verifierMonitoring, secrets) + chainStatusManager, chainStatusDB, err := createChainStatusManager(ctx, lggr, config.VerifierID, verifierMonitoring, secrets) if err != nil { lggr.Errorw("Failed to create chain status manager", "error", err) return fmt.Errorf("failed to create chain status manager: %w", err) @@ -542,8 +563,8 @@ func (f *factory) Stop(ctx context.Context) error { return allErrors } -func createChainStatusManager(lggr logger.Logger, verifierID string, monitoring verifier.Monitoring, secrets *vsecrets.VerifierSecrets) (protocol.ChainStatusManager, sqlutil.DataSource, error) { - sqlDB, err := ConnectToPostgresDB(lggr, secrets) +func createChainStatusManager(ctx context.Context, lggr logger.Logger, verifierID string, monitoring verifier.Monitoring, secrets *vsecrets.VerifierSecrets) (protocol.ChainStatusManager, sqlutil.DataSource, error) { + sqlDB, err := ConnectToPostgresDB(ctx, lggr, secrets) if err != nil { return nil, nil, fmt.Errorf("failed to connect to Postgres DB: %w", err) } diff --git a/cmd/verifier/tokenfactory.go b/cmd/verifier/tokenfactory.go index 1e71f8044..9c873db32 100644 --- a/cmd/verifier/tokenfactory.go +++ b/cmd/verifier/tokenfactory.go @@ -7,9 +7,11 @@ import ( "net/http" "slices" "strconv" + "sync" "time" sdkmetric "go.opentelemetry.io/otel/sdk/metric" + "golang.org/x/sync/errgroup" "github.com/smartcontractkit/chainlink-common/pkg/sqlutil" @@ -99,25 +101,40 @@ func (tvf *tokenVerifierFactory) Start(ctx context.Context, spec bootstrap.JobSp cfg := appConfig // On-ramp addresses are the application-owned source-chain set. RPC connection and tuning - // details are supplied independently by each chain family's local config. + // details are supplied independently by each chain family's local config. Chains are built + // concurrently with a per-chain timeout: accessor construction dials the chain, so a slow + // (not failing) RPC must not serialize away the shared startup budget. chainSelectors := chainaccess.Infos[string](cfg.OnRampAddresses).GetAllChainSelectors() sourceReaders := make(map[protocol.ChainSelector]chainaccess.SourceReader) accessors := make(map[protocol.ChainSelector]chainaccess.Accessor) + var accessorsMu sync.Mutex + g, gctx := errgroup.WithContext(ctx) for _, selector := range chainSelectors { - accessor, err := deps.Registry.GetAccessor(ctx, selector) - if err != nil { - tvf.lggr.Errorw("Skipping chain, failed to get accessor for chain selector", "error", err, "chainSelector", selector) - continue - } - accessors[selector] = accessor - reader, err := accessor.SourceReader() - if err != nil { - tvf.lggr.Warnw("Skipping chain, source reader not available", "chainSelector", selector, "error", err) - continue - } - sourceReaders[selector] = reader - tvf.lggr.Infow("Created source reader for chain", "chainSelector", selector) + g.Go(func() error { + chainCtx, cancel := context.WithTimeout(gctx, accessorBuildTimeout) + defer cancel() + //nolint:noeagerio // accessor construction dials the chain; it runs concurrently per chain with a bounded ctx, and failures are skipped, not fatal + accessor, err := deps.Registry.GetAccessor(chainCtx, selector) + if err != nil { + tvf.lggr.Errorw("Skipping chain, failed to get accessor for chain selector", "error", err, "chainSelector", selector) + return nil + } + reader, err := accessor.SourceReader() + accessorsMu.Lock() + accessors[selector] = accessor + accessorsMu.Unlock() + if err != nil { + tvf.lggr.Warnw("Skipping chain, source reader not available", "chainSelector", selector, "error", err) + return nil + } + accessorsMu.Lock() + sourceReaders[selector] = reader + accessorsMu.Unlock() + tvf.lggr.Infow("Created source reader for chain", "chainSelector", selector) + return nil + }) } + _ = g.Wait() // Load the verifier secrets file (only [db].url is used by the token verifier); an absent file // is fine and falls back to CL_DATABASE_URL. @@ -126,7 +143,7 @@ func (tvf *tokenVerifierFactory) Start(ctx context.Context, spec bootstrap.JobSp return fmt.Errorf("failed to load verifier secrets: %w", err) } - db, err := ConnectToPostgresDB(tvf.lggr, secrets) + db, err := ConnectToPostgresDB(ctx, tvf.lggr, secrets) if err != nil { return fmt.Errorf("failed to connect to Postgres database: %w", err) } diff --git a/common/db_utils.go b/common/db_utils.go index d90369735..36cb3e93e 100644 --- a/common/db_utils.go +++ b/common/db_utils.go @@ -21,22 +21,37 @@ type Pingable interface { } // EnsureDBConnection ensures that the database is up and running by pinging it. +// +// Deprecated: use EnsureDBConnectionContext so a caller's startup deadline can +// cancel the retries; this wrapper is unbounded by any caller context. func EnsureDBConnection(lggr logger.Logger, db Pingable) error { - pingFn := func() error { - ctx, cancel := context.WithTimeout(context.Background(), timeout) - defer cancel() - return db.PingContext(ctx) - } - for range maxRetries { - err := pingFn() + return EnsureDBConnectionContext(context.Background(), lggr, db) +} + +// EnsureDBConnectionContext pings the database until it answers or ctx is done. +// Retries stop as soon as ctx is canceled, so a degraded database cannot block +// startup beyond the caller's deadline. +func EnsureDBConnectionContext(ctx context.Context, lggr logger.Logger, db Pingable) error { + for attempt := range maxRetries { + pingCtx, cancel := context.WithTimeout(ctx, timeout) + err := db.PingContext(pingCtx) + cancel() if err == nil { return nil } + if ctx.Err() != nil { + return fmt.Errorf("database still unreachable (last ping: %w): %w", err, ctx.Err()) + } lggr.Warnw("failed to connect to database, retrying after sleeping", "err", err, "retryInterval", retryInterval.String(), + "attempt", attempt+1, "maxRetries", maxRetries) - time.Sleep(retryInterval) + select { + case <-ctx.Done(): + return fmt.Errorf("database still unreachable (last ping: %w): %w", err, ctx.Err()) + case <-time.After(retryInterval): + } } return fmt.Errorf("failed to connect to database after %d retries", maxRetries) } diff --git a/common/jobqueue/postgres_queue_explain_test.go b/common/jobqueue/postgres_queue_explain_test.go index 2fa03b2cc..55246ff26 100644 --- a/common/jobqueue/postgres_queue_explain_test.go +++ b/common/jobqueue/postgres_queue_explain_test.go @@ -55,7 +55,7 @@ func getExplainDB(t *testing.T) *sqlx.DB { if url := os.Getenv("TEST_POSTGRES_URL"); url != "" { sdb, err := sqlx.Open("postgres", url) require.NoError(t, err, "open local postgres") - require.NoError(t, verifierdb.RunPostgresMigrations(sdb), "run migrations") + require.NoError(t, verifierdb.RunPostgresMigrationsContext(t.Context(), sdb), "run migrations") _, err = sdb.ExecContext(context.Background(), "TRUNCATE ccv_task_verifier_jobs, ccv_task_verifier_jobs_archive CASCADE") require.NoError(t, err, "truncate tables") diff --git a/indexer/cmd/main.go b/indexer/cmd/main.go index 3cb0c45ff..f780eb177 100644 --- a/indexer/cmd/main.go +++ b/indexer/cmd/main.go @@ -323,10 +323,10 @@ func createPostgresStorage(ctx context.Context, lggr logger.Logger, cfg *config. if err != nil { lggr.Fatalf("Failed to open database for migrations: %v", err) } - if err := ccvcommon.EnsureDBConnection(lggr, migrationsDB.DB); err != nil { + if err := ccvcommon.EnsureDBConnectionContext(ctx, lggr, migrationsDB.DB); err != nil { lggr.Fatalf("Could not connect to database: %v", err) } - if err := storage.RunMigrations(migrationsDB); err != nil { + if err := storage.RunMigrationsContext(ctx, migrationsDB); err != nil { lggr.Fatalf("Failed to run database migrations: %v", err) } if err := migrationsDB.Close(); err != nil { diff --git a/indexer/cmd/replay/main.go b/indexer/cmd/replay/main.go index 9a7146980..a16c2c375 100644 --- a/indexer/cmd/replay/main.go +++ b/indexer/cmd/replay/main.go @@ -314,10 +314,10 @@ func mustBuildEngine(ctx context.Context, needsDiscoveryReader bool) (*replay.En if err != nil { lggr.Fatalf("Failed to open database for migrations: %v", err) } - if err := ccvcommon.EnsureDBConnection(lggr, migrationsDB.DB); err != nil { + if err := ccvcommon.EnsureDBConnectionContext(ctx, lggr, migrationsDB.DB); err != nil { lggr.Fatalf("Could not connect to database: %v", err) } - if err := storage.RunMigrations(migrationsDB); err != nil { + if err := storage.RunMigrationsContext(ctx, migrationsDB); err != nil { lggr.Fatalf("Failed to run database migrations: %v", err) } if err := migrationsDB.Close(); err != nil { diff --git a/indexer/pkg/storage/run_migrations.go b/indexer/pkg/storage/run_migrations.go index f5636ecc9..a6dc5f8fe 100644 --- a/indexer/pkg/storage/run_migrations.go +++ b/indexer/pkg/storage/run_migrations.go @@ -1,6 +1,7 @@ package storage import ( + "context" "fmt" "sync" @@ -13,7 +14,16 @@ import ( var migrationMutex = sync.Mutex{} // RunMigrations applies PostgreSQL database migrations. +// +// Deprecated: use RunMigrationsContext so a caller's startup deadline can abort +// a hung migration; this wrapper is unbounded by any caller context. func RunMigrations(db *sqlx.DB) error { + return RunMigrationsContext(context.Background(), db) +} + +// RunMigrationsContext applies PostgreSQL database migrations, aborting when +// ctx is done. +func RunMigrationsContext(ctx context.Context, db *sqlx.DB) error { migrationMutex.Lock() defer migrationMutex.Unlock() @@ -23,7 +33,7 @@ func RunMigrations(db *sqlx.DB) error { return fmt.Errorf("failed to set goose dialect: %w", err) } - if err := goose.Up(db.DB, "postgres"); err != nil { + if err := goose.UpContext(ctx, db.DB, "postgres"); err != nil { return fmt.Errorf("failed to run postgres migrations: %w", err) } diff --git a/integration/pkg/messagerules/poller.go b/integration/pkg/messagerules/poller.go index 5da57ef96..686172ec0 100644 --- a/integration/pkg/messagerules/poller.go +++ b/integration/pkg/messagerules/poller.go @@ -61,7 +61,12 @@ func NewPollerService(client Client, pollInterval, clientTimeout time.Duration, func (s *PollerService) Start(ctx context.Context) error { return s.StartOnce(s.Name(), func() error { s.wg.Go(func() { - s.poll(ctx) + // The initial poll must use the service-lifetime context, not the + // startup ctx: bootstrap cancels the startup ctx as soon as Start + // returns, which could abort the very first rules fetch. + pollCtx, cancel := s.stopCh.NewCtx() + defer cancel() + s.poll(pollCtx) s.pollLoop() }) diff --git a/tools/noeagerio/analyzer.go b/tools/noeagerio/analyzer.go index 591685a11..0d306b6ae 100644 --- a/tools/noeagerio/analyzer.go +++ b/tools/noeagerio/analyzer.go @@ -99,6 +99,7 @@ var denyMethods = map[string]string{ "LoadKeystore": "reads from the keystore", "GetPublicKey": "queries the KMS", "LoadKMSKeystore": "queries the KMS", + "GetAccessor": "constructs a chain accessor (dials RPC, starts chain services)", // Known repo entrypoints whose I/O happens across a package boundary and so // cannot be found by local taint propagation. "New@sqlutil/pg": "opens a database connection", diff --git a/tools/noeagerio/testdata/src/a/a.go b/tools/noeagerio/testdata/src/a/a.go index bf8047cd9..bdfedd2d4 100644 --- a/tools/noeagerio/testdata/src/a/a.go +++ b/tools/noeagerio/testdata/src/a/a.go @@ -97,3 +97,20 @@ func NewLazyService(d *dialer) *lazyService { }), } } + +// registry stands in for chainaccess registries: GetAccessor is denylisted by +// bare method name because accessor construction dials the chain. +type registry struct{} + +func (r *registry) GetAccessor(ctx context.Context, selector uint64) (*dialer, error) { + return &dialer{}, nil +} + +type factoryService struct { + reg *registry +} + +func (s *factoryService) Start(ctx context.Context) error { + _, err := s.reg.GetAccessor(ctx, 1) // want `I/O in Start method Start: constructs a chain accessor` + return err +} diff --git a/verifier/pkg/coordinator.go b/verifier/pkg/coordinator.go index 140ca8a82..88cd21631 100644 --- a/verifier/pkg/coordinator.go +++ b/verifier/pkg/coordinator.go @@ -41,6 +41,8 @@ const ( resultQueueLockDuration = 1 * time.Minute // queueObservabilityInterval is how often queue size metrics are logged and recorded. queueObservabilityInterval = 10 * time.Second + // chainStatusReadTimeout bounds the best-effort chain status read at coordinator start. + chainStatusReadTimeout = 5 * time.Second ) type Coordinator struct { @@ -448,10 +450,12 @@ func filterConfiguredSourceReaders( } // A transient DB failure here must not abort startup: the statuses only - // inform logging and the disabled-chain warning below. Degrade to unknown - // statuses; each source reader re-reads its own status (with retries) when - // it initializes its start block. - statusMap, err := chainStatusManager.ReadChainStatuses(ctx, allSelectors) //nolint:noeagerio // single bounded read of the service's own DB at coordinator start; failure is non-fatal + // inform logging and the disabled-chain warning below. Bound the read and + // degrade to unknown statuses; each source reader re-reads its own status + // (with retries) when it initializes its start block. + readCtx, cancel := context.WithTimeout(ctx, chainStatusReadTimeout) + statusMap, err := chainStatusManager.ReadChainStatuses(readCtx, allSelectors) //nolint:noeagerio // single bounded read of the service's own DB at coordinator start; failure is non-fatal + cancel() if err != nil { lggr.Errorw("Failed to read chain statuses from storage, continuing with unknown statuses", "error", err) statusMap = nil diff --git a/verifier/pkg/db/run_migrations.go b/verifier/pkg/db/run_migrations.go index 9acb46267..5edf08e6a 100644 --- a/verifier/pkg/db/run_migrations.go +++ b/verifier/pkg/db/run_migrations.go @@ -1,6 +1,7 @@ package db import ( + "context" "fmt" "sync" @@ -13,7 +14,16 @@ import ( var migrationMutex = sync.Mutex{} // RunPostgresMigrations applies PostgreSQL database migrations. +// +// Deprecated: use RunPostgresMigrationsContext so a caller's startup deadline +// can abort a hung migration; this wrapper is unbounded by any caller context. func RunPostgresMigrations(db *sqlx.DB) error { + return RunPostgresMigrationsContext(context.Background(), db) +} + +// RunPostgresMigrationsContext applies PostgreSQL database migrations, aborting +// when ctx is done. +func RunPostgresMigrationsContext(ctx context.Context, db *sqlx.DB) error { migrationMutex.Lock() defer migrationMutex.Unlock() @@ -23,7 +33,7 @@ func RunPostgresMigrations(db *sqlx.DB) error { return fmt.Errorf("failed to set goose dialect: %w", err) } - if err := goose.Up(db.DB, "postgres"); err != nil { + if err := goose.UpContext(ctx, db.DB, "postgres"); err != nil { return fmt.Errorf("failed to run postgres migrations: %w", err) } diff --git a/verifier/testutil/test_db.go b/verifier/testutil/test_db.go index 806cd791f..799589bc8 100644 --- a/verifier/testutil/test_db.go +++ b/verifier/testutil/test_db.go @@ -45,7 +45,7 @@ func NewTestDB(tb testing.TB) *sqlx.DB { sqlxDB, err := sqlx.Open("postgres", connectionString) require.NoError(tb, err, "failed to open database") - err = db.RunPostgresMigrations(sqlxDB) + err = db.RunPostgresMigrationsContext(ctx, sqlxDB) require.NoError(tb, err, "failed to run migrations") tb.Cleanup(func() {