Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions .github/workflows/golangci-lint.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
13 changes: 13 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
6 changes: 6 additions & 0 deletions Justfile
Original file line number Diff line number Diff line change
Expand Up @@ -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."; \
Expand Down
6 changes: 5 additions & 1 deletion aggregator/cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -210,11 +210,15 @@

protocol.InitChainSelectorCache()

server := aggregator.NewServer(sugaredLggr, config, aggMonitoring)
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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
19 changes: 11 additions & 8 deletions aggregator/pkg/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -309,9 +309,12 @@ 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. 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 {
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",
Expand All @@ -325,10 +328,10 @@ func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonito
)

factory := storage.NewStorageFactory(l)
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 {
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.
Expand Down Expand Up @@ -365,14 +368,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 {
Expand Down Expand Up @@ -470,5 +473,5 @@ func NewServer(l logger.SugaredLogger, config *model.AggregatorConfig, aggMonito
committeepb.RegisterCommitteeVerifierServer(grpcServer, server)
heartbeatpb.RegisterHeartbeatServiceServer(grpcServer, server)

return server
return server, nil
}
12 changes: 7 additions & 5 deletions aggregator/pkg/storage/factory.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package storage

import (
"context"
"database/sql"
"fmt"
"time"
Expand Down Expand Up @@ -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")
}
Expand Down Expand Up @@ -98,15 +100,15 @@ 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)
}

// Create sqlx wrapper for sqlutil.DataSource compatibility
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)
}
Expand Down
12 changes: 11 additions & 1 deletion aggregator/pkg/storage/postgres/run_migrations.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package postgres

import (
"context"
"fmt"
"sync"

Expand All @@ -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()

Expand All @@ -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)
}

Expand Down
8 changes: 6 additions & 2 deletions aggregator/tests/utils.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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(t.Context(), 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)
Expand Down
4 changes: 3 additions & 1 deletion bootstrap/bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}

Expand Down Expand Up @@ -772,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
Expand Down
2 changes: 1 addition & 1 deletion bootstrap/bootstrap_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 12 additions & 1 deletion bootstrap/db/run_migrations.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package db

import (
"context"
"embed"
"fmt"
"sync"
Expand All @@ -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()

Expand All @@ -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)
}

Expand Down
1 change: 1 addition & 0 deletions bootstrap/keys/csa.go
Original file line number Diff line number Diff line change
Expand Up @@ -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},
})
Expand Down
1 change: 1 addition & 0 deletions bootstrap/keys/kms.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
55 changes: 55 additions & 0 deletions changelog/2026-10-06_no_eager_io_startup.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
# 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 `(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
read, the aggregator's own storage connect+migrate, the opt-in Redis rate
limiter, and the operator-run indexer replay tool.
Loading
Loading