diff --git a/docs/config/indexer/config.documented.toml b/docs/config/indexer/config.documented.toml index f3aea4dc2..cc956be09 100644 --- a/docs/config/indexer/config.documented.toml +++ b/docs/config/indexer/config.documented.toml @@ -166,3 +166,37 @@ MergeBufferSize = 0 # Enabled enables the rate limiting system inside the indexer. Enabled = false +# Resilience configures the resilient reader policies (rate limiting, bulkheading, +# circuit breaker, and request timeouts). +[Resilience] + # MaxRequestsPerSecond is the maximum number of requests per second allowed per reader. + # 0 uses the default (5). + MaxRequestsPerSecond = 0 + # MaxConcurrentRequests is the maximum number of concurrent requests allowed per reader. + # 0 uses the default (5). + MaxConcurrentRequests = 0 + # FailureThreshold is the number of consecutive failures that opens the circuit breaker. + # 0 uses the default (5). + FailureThreshold = 0 + # SuccessThreshold is the number of consecutive successes that closes the circuit breaker. + # 0 uses the default (3). + SuccessThreshold = 0 + # CircuitBreakerDelay is how long the circuit breaker stays open before entering half-open. + # 0 uses the default (3s). + CircuitBreakerDelay = "0s" + # CircuitBreakerTimeout is the timeout for circuit breaker operations. + # 0 uses the default (1s). + CircuitBreakerTimeout = "0s" + # RequestTimeout is the per-request timeout. + # 0 uses the default (10s). + RequestTimeout = "0s" + # MaxRetries is the maximum number of retry attempts per request. + # 0 uses the default (3). + MaxRetries = 0 + # RetryDelay is the initial delay between retries, using exponential backoff. + # 0 uses the default (1s). + RetryDelay = "0s" + # RetryMaxDelay is the maximum delay between retries. + # 0 uses the default (10s). + RetryMaxDelay = "0s" + diff --git a/indexer/cmd/main.go b/indexer/cmd/main.go index eadaee451..bd76c837f 100644 --- a/indexer/cmd/main.go +++ b/indexer/cmd/main.go @@ -145,7 +145,7 @@ func createRegistry() *registry.VerifierRegistry { func createAllVerifierReaders(ctx context.Context, lggr logger.Logger, verifierRegistry *registry.VerifierRegistry, config *config.Config, indexerMonitoring common.IndexerMonitoring) error { for _, verifierConfig := range config.Verifiers { - err := createReadersForVerifier(ctx, lggr, verifierRegistry, &verifierConfig, indexerMonitoring) + err := createReadersForVerifier(ctx, lggr, verifierRegistry, &verifierConfig, indexerMonitoring, config.Resilience) if err != nil { return err } @@ -154,9 +154,9 @@ func createAllVerifierReaders(ctx context.Context, lggr logger.Logger, verifierR return nil } -func createReadersForVerifier(ctx context.Context, lggr logger.Logger, verifierRegistry *registry.VerifierRegistry, verifierConfig *config.VerifierConfig, monitoring common.IndexerMonitoring) error { +func createReadersForVerifier(ctx context.Context, lggr logger.Logger, verifierRegistry *registry.VerifierRegistry, verifierConfig *config.VerifierConfig, monitoring common.IndexerMonitoring, resilience config.ResilienceConfig) error { metrics := monitoring.Metrics().With("target", verifierConfig.Name) - reader, err := createReader(lggr, verifierConfig, metrics) + reader, err := createReader(lggr, verifierConfig, metrics, resilience) if err != nil { return err } @@ -182,13 +182,13 @@ func createReadersForVerifier(ctx context.Context, lggr logger.Logger, verifierR return nil } -func createReader(lggr logger.Logger, cfg *config.VerifierConfig, m common.IndexerMetricLabeler) (*readers.ResilientReader, error) { +func createReader(lggr logger.Logger, cfg *config.VerifierConfig, m common.IndexerMetricLabeler, resiConfig config.ResilienceConfig) (*readers.ResilientReader, error) { switch cfg.Type { case config.ReaderTypeAggregator: return readers.NewAggregatorReader(cfg.Address, lggr, cfg.Since, hmac.ClientConfig{ APIKey: cfg.APIKey, Secret: cfg.Secret, - }, cfg.InsecureConnection, config.EffectiveMaxResponseBytes(cfg.MaxResponseBytes), m) + }, cfg.InsecureConnection, config.EffectiveMaxResponseBytes(cfg.MaxResponseBytes), m, readers.NewResilienceConfig(resiConfig)) case config.ReaderTypeRest: return readers.NewRestReader(readers.RestReaderConfig{ BaseURL: cfg.BaseURL, @@ -196,6 +196,7 @@ func createReader(lggr logger.Logger, cfg *config.VerifierConfig, m common.Index MaxResponseBytes: config.EffectiveMaxResponseBytes(cfg.MaxResponseBytes), Logger: lggr, Metrics: m, + Resilience: readers.NewResilienceConfig(resiConfig), }), nil default: return nil, errors.New("unknown verifier type") @@ -235,7 +236,7 @@ func createDiscovery(ctx context.Context, lggr logger.Logger, cfg *config.Config aggregator, err := readers.NewAggregatorReader(discCfg.Address, lggr, int64(persistedSinceValue), hmac.ClientConfig{ APIKey: discCfg.APIKey, Secret: discCfg.Secret, - }, discCfg.InsecureConnection, config.EffectiveMaxResponseBytes(discCfg.MaxResponseBytes), metrics) + }, discCfg.InsecureConnection, config.EffectiveMaxResponseBytes(discCfg.MaxResponseBytes), metrics, readers.NewResilienceConfig(cfg.Resilience)) if err != nil { cleanupOnError() return nil, err diff --git a/indexer/cmd/replay/main.go b/indexer/cmd/replay/main.go index f224f3ce9..9a7146980 100644 --- a/indexer/cmd/replay/main.go +++ b/indexer/cmd/replay/main.go @@ -337,7 +337,7 @@ func mustBuildEngine(ctx context.Context, needsDiscoveryReader bool) (*replay.En verifierRegistry := registry.NewVerifierRegistry() verifierCleanups := make([]func(), 0, len(cfg.Verifiers)) for _, vc := range cfg.Verifiers { - vr, cleanup, err := createVerifierReader(ctx, lggr, &vc, monitoring) + vr, cleanup, err := createVerifierReader(ctx, lggr, &vc, monitoring, cfg.Resilience) if err != nil { lggr.Fatalf("Failed to create verifier reader: %v", err) } @@ -362,7 +362,7 @@ func mustBuildEngine(ctx context.Context, needsDiscoveryReader bool) (*replay.En return readers.NewAggregatorReader(disc.Address, lggr, since, hmac.ClientConfig{ APIKey: disc.APIKey, Secret: disc.Secret, - }, disc.InsecureConnection, config.EffectiveMaxResponseBytes(disc.MaxResponseBytes), metrics) + }, disc.InsecureConnection, config.EffectiveMaxResponseBytes(disc.MaxResponseBytes), metrics, readers.NewResilienceConfig(cfg.Resilience)) } } @@ -418,7 +418,7 @@ func mustCreateLogger(cfg *config.Config) logger.Logger { return logger.Named(logger.Sugared(lggr), "indexer-replay") } -func createVerifierReader(ctx context.Context, lggr logger.Logger, vc *config.VerifierConfig, mon common.IndexerMonitoring) (*readers.VerifierReader, func(), error) { +func createVerifierReader(ctx context.Context, lggr logger.Logger, vc *config.VerifierConfig, mon common.IndexerMonitoring, resilience config.ResilienceConfig) (*readers.VerifierReader, func(), error) { var resilientReader *readers.ResilientReader var err error @@ -428,7 +428,7 @@ func createVerifierReader(ctx context.Context, lggr logger.Logger, vc *config.Ve resilientReader, err = readers.NewAggregatorReader(vc.Address, lggr, vc.Since, hmac.ClientConfig{ APIKey: vc.APIKey, Secret: vc.Secret, - }, vc.InsecureConnection, config.EffectiveMaxResponseBytes(vc.MaxResponseBytes), metrics) + }, vc.InsecureConnection, config.EffectiveMaxResponseBytes(vc.MaxResponseBytes), metrics, readers.NewResilienceConfig(resilience)) case config.ReaderTypeRest: resilientReader = readers.NewRestReader(readers.RestReaderConfig{ BaseURL: vc.BaseURL, @@ -436,6 +436,7 @@ func createVerifierReader(ctx context.Context, lggr logger.Logger, vc *config.Ve MaxResponseBytes: config.EffectiveMaxResponseBytes(vc.MaxResponseBytes), Logger: lggr, Metrics: metrics, + Resilience: readers.NewResilienceConfig(resilience), }) default: return nil, nil, errors.New("unknown verifier type: " + string(vc.Type)) diff --git a/indexer/config.example.toml b/indexer/config.example.toml index 7c92f465b..c0e0a720e 100644 --- a/indexer/config.example.toml +++ b/indexer/config.example.toml @@ -52,3 +52,15 @@ LockTimeout = 30 [API] [API.RateLimit] Enabled = false + +[Resilience] +MaxRequestsPerSecond = 5 +MaxConcurrentRequests = 5 +FailureThreshold = 5 +SuccessThreshold = 3 +CircuitBreakerDelay = "3s" +CircuitBreakerTimeout = "1s" +RequestTimeout = "10s" +MaxRetries = 3 +RetryDelay = "1s" +RetryMaxDelay = "10s" diff --git a/indexer/pkg/config/config.go b/indexer/pkg/config/config.go index c72b05f56..84e9e5266 100644 --- a/indexer/pkg/config/config.go +++ b/indexer/pkg/config/config.go @@ -51,6 +51,44 @@ type Config struct { Storage StorageConfig `toml:"Storage"` // API is the configuration for the API inside the indexer. API APIConfig `toml:"API"` + // Resilience configures the resilient reader policies (rate limiting, bulkheading, + // circuit breaker, and request timeouts). + Resilience ResilienceConfig `toml:"Resilience"` +} + +// ResilienceConfig provides configuration for the resilient reader policies. +// Zero values fall back to the defaults in readers.DefaultResilienceConfig(). +type ResilienceConfig struct { + // MaxRequestsPerSecond is the maximum number of requests per second allowed per reader. + // 0 uses the default (5). + MaxRequestsPerSecond uint `toml:"MaxRequestsPerSecond"` + // MaxConcurrentRequests is the maximum number of concurrent requests allowed per reader. + // 0 uses the default (5). + MaxConcurrentRequests uint `toml:"MaxConcurrentRequests"` + // FailureThreshold is the number of consecutive failures that opens the circuit breaker. + // 0 uses the default (5). + FailureThreshold uint32 `toml:"FailureThreshold"` + // SuccessThreshold is the number of consecutive successes that closes the circuit breaker. + // 0 uses the default (3). + SuccessThreshold uint32 `toml:"SuccessThreshold"` + // CircuitBreakerDelay is how long the circuit breaker stays open before entering half-open. + // 0 uses the default (3s). + CircuitBreakerDelay common.Duration `toml:"CircuitBreakerDelay"` + // CircuitBreakerTimeout is the timeout for circuit breaker operations. + // 0 uses the default (1s). + CircuitBreakerTimeout common.Duration `toml:"CircuitBreakerTimeout"` + // RequestTimeout is the per-request timeout. + // 0 uses the default (10s). + RequestTimeout common.Duration `toml:"RequestTimeout"` + // MaxRetries is the maximum number of retry attempts per request. + // 0 uses the default (3). + MaxRetries int `toml:"MaxRetries"` + // RetryDelay is the initial delay between retries, using exponential backoff. + // 0 uses the default (1s). + RetryDelay common.Duration `toml:"RetryDelay"` + // RetryMaxDelay is the maximum delay between retries. + // 0 uses the default (10s). + RetryMaxDelay common.Duration `toml:"RetryMaxDelay"` } type SchedulerConfig struct { diff --git a/indexer/pkg/readers/aggregator_reader.go b/indexer/pkg/readers/aggregator_reader.go index 6c64f5b44..c1a4a986e 100644 --- a/indexer/pkg/readers/aggregator_reader.go +++ b/indexer/pkg/readers/aggregator_reader.go @@ -11,18 +11,17 @@ import ( "github.com/smartcontractkit/chainlink-common/pkg/logger" ) -func NewAggregatorReader(address string, lggr logger.Logger, since int64, hmacConfig hmac.ClientConfig, insecure bool, maxRecvMsgSizeBytes int, m common.IndexerMetricLabeler) (*ResilientReader, error) { +func NewAggregatorReader(address string, lggr logger.Logger, since int64, hmacConfig hmac.ClientConfig, insecure bool, maxRecvMsgSizeBytes int, m common.IndexerMetricLabeler, resiConfig ResilienceConfig) (*ResilientReader, error) { reader, err := storageaccess.NewAggregatorReader(address, lggr, since, &hmacConfig, insecure, maxRecvMsgSizeBytes, grpcClientDialOptions(m)...) if err != nil { return nil, fmt.Errorf("failed to create aggregator reader: %w", err) } observed := NewObservedReader(reader, reader, m) - config := DefaultResilienceConfig() - config.DiscoveryRetryPolicyErrorHandler = aggregatorRetryPolicyErrorHandler - config.DiscoveryCircuitBreakerErrorHandler = aggregatorRetryPolicyErrorHandler + resiConfig.DiscoveryRetryPolicyErrorHandler = aggregatorRetryPolicyErrorHandler + resiConfig.DiscoveryCircuitBreakerErrorHandler = aggregatorRetryPolicyErrorHandler - return NewResilientReader(observed, lggr, config), nil + return NewResilientReader(observed, lggr, resiConfig), nil } // aggregatorRetryPolicyErrorHandler determines if an error from the aggregator should be retried. diff --git a/indexer/pkg/readers/resilient_reader.go b/indexer/pkg/readers/resilient_reader.go index 9cf6332ea..56179d878 100644 --- a/indexer/pkg/readers/resilient_reader.go +++ b/indexer/pkg/readers/resilient_reader.go @@ -10,8 +10,10 @@ import ( "github.com/failsafe-go/failsafe-go/bulkhead" "github.com/failsafe-go/failsafe-go/circuitbreaker" "github.com/failsafe-go/failsafe-go/ratelimiter" + "github.com/failsafe-go/failsafe-go/retrypolicy" "github.com/failsafe-go/failsafe-go/timeout" + "github.com/smartcontractkit/chainlink-ccv/indexer/pkg/config" "github.com/smartcontractkit/chainlink-ccv/protocol" "github.com/smartcontractkit/chainlink-common/pkg/logger" ) @@ -37,6 +39,9 @@ type ResilienceConfig struct { RequestTimeout time.Duration MaxConcurrentRequests uint MaxRequestsPerSecond uint + MaxRetries int + RetryDelay time.Duration + RetryMaxDelay time.Duration } // DefaultResilienceConfig returns a configuration with sensible defaults. @@ -49,9 +54,49 @@ func DefaultResilienceConfig() ResilienceConfig { RequestTimeout: 10 * time.Second, MaxConcurrentRequests: 5, MaxRequestsPerSecond: 5, + MaxRetries: 3, + RetryDelay: 1 * time.Second, + RetryMaxDelay: 10 * time.Second, } } +// NewResilienceConfig builds a readers.ResilienceConfig from the indexer's +// config.ResilienceConfig, applying defaults for any zero-value fields. +func NewResilienceConfig(c config.ResilienceConfig) ResilienceConfig { + rc := DefaultResilienceConfig() + if c.MaxRequestsPerSecond > 0 { + rc.MaxRequestsPerSecond = c.MaxRequestsPerSecond + } + if c.MaxConcurrentRequests > 0 { + rc.MaxConcurrentRequests = c.MaxConcurrentRequests + } + if c.FailureThreshold > 0 { + rc.FailureThreshold = c.FailureThreshold + } + if c.SuccessThreshold > 0 { + rc.SuccessThreshold = c.SuccessThreshold + } + if c.CircuitBreakerDelay > 0 { + rc.CircuitBreakerDelay = time.Duration(c.CircuitBreakerDelay) + } + if c.CircuitBreakerTimeout > 0 { + rc.CircuitBreakerTimeout = time.Duration(c.CircuitBreakerTimeout) + } + if c.RequestTimeout > 0 { + rc.RequestTimeout = time.Duration(c.RequestTimeout) + } + if c.MaxRetries > 0 { + rc.MaxRetries = c.MaxRetries + } + if c.RetryDelay > 0 { + rc.RetryDelay = time.Duration(c.RetryDelay) + } + if c.RetryMaxDelay > 0 { + rc.RetryMaxDelay = time.Duration(c.RetryMaxDelay) + } + return rc +} + type executorPolicies[T any] struct { executor failsafe.Executor[T] circuitBreaker circuitbreaker.CircuitBreaker[T] @@ -78,25 +123,41 @@ func NewResilientReader(underlying protocol.VerifierResultsAPI, lggr logger.Logg maxConsecutiveErrors: config.FailureThreshold, } - rr.verificationsPolicies = createPolicies(config, lggr, "GetVerifications", config.CircuitBreakerErrorHandler) + rr.verificationsPolicies = createPolicies(config, lggr, "GetVerifications", config.RetryPolicyErrorHandler, config.CircuitBreakerErrorHandler) if discoveryAPI, ok := underlying.(protocol.OffchainStorageReader); ok { - rr.discoveryPolicies = createPolicies(config, lggr, "ReadCCVData", config.DiscoveryCircuitBreakerErrorHandler) + rr.discoveryPolicies = createPolicies(config, lggr, "ReadCCVData", config.DiscoveryRetryPolicyErrorHandler, config.DiscoveryCircuitBreakerErrorHandler) rr.discoveryAPI = discoveryAPI } return rr } -func createPolicies[T any](config ResilienceConfig, lggr logger.Logger, name string, errorHandler func(T, error) bool) executorPolicies[T] { - handleIf := func(resp T, err error) bool { return err != nil } - if errorHandler != nil { - handleIf = errorHandler +func createPolicies[T any](config ResilienceConfig, lggr logger.Logger, name string, retryErrorHandler, cbErrorHandler func(T, error) bool) executorPolicies[T] { + retryHandleIf := func(resp T, err error) bool { return err != nil } + if retryErrorHandler != nil { + retryHandleIf = retryErrorHandler + } + + rp := retrypolicy.NewBuilder[T](). + HandleIf(retryHandleIf). + WithMaxRetries(config.MaxRetries). + WithBackoff(config.RetryDelay, config.RetryMaxDelay). + AbortOnErrors(context.Canceled, context.DeadlineExceeded, circuitbreaker.ErrOpen). + ReturnLastFailure(). + OnRetry(func(failsafe.ExecutionEvent[T]) { + lggr.Warnw(name+" retrying request", "max_retries", config.MaxRetries) + }). + Build() + + cbHandleIf := func(resp T, err error) bool { return err != nil } + if cbErrorHandler != nil { + cbHandleIf = cbErrorHandler } cb := circuitbreaker.NewBuilder[T](). WithDelay(config.CircuitBreakerDelay). - HandleIf(handleIf). + HandleIf(cbHandleIf). OnOpen(func(circuitbreaker.StateChangedEvent) { lggr.Warnw(name+" circuit breaker opened", "failures", config.FailureThreshold) }). @@ -110,7 +171,9 @@ func createPolicies[T any](config ResilienceConfig, lggr logger.Logger, name str WithSuccessThreshold(uint(config.SuccessThreshold)). Build() - rl := ratelimiter.NewBursty[T](config.MaxRequestsPerSecond, time.Second) + rl := ratelimiter.NewBurstyBuilder[T](config.MaxRequestsPerSecond, time.Second). + WithMaxWaitTime(time.Second). // Wait up to 1 second for a permit before returning ErrExceeded + Build() bh := bulkhead.NewBuilder[T](config.MaxConcurrentRequests). OnFull(func(failsafe.ExecutionEvent[T]) { lggr.Warnw(name+" bulkhead is full", "max_concurrent_requests", config.MaxConcurrentRequests) @@ -123,25 +186,25 @@ func createPolicies[T any](config ResilienceConfig, lggr logger.Logger, name str Build() return executorPolicies[T]{ - executor: failsafe.With(cb, rl, bh, to), + executor: failsafe.With(rp, cb, rl, bh, to), circuitBreaker: cb, } } func (r *ResilientReader) ReadCCVData(ctx context.Context) ([]protocol.QueryResponse, error) { - return execute(r, r.discoveryPolicies, func() ([]protocol.QueryResponse, error) { + return execute(ctx, r, r.discoveryPolicies, func() ([]protocol.QueryResponse, error) { return r.discoveryAPI.ReadCCVData(ctx) }) } func (r *ResilientReader) GetVerifications(ctx context.Context, messageIDs []protocol.Bytes32) (map[protocol.Bytes32]protocol.VerifierResult, error) { - return execute(r, r.verificationsPolicies, func() (map[protocol.Bytes32]protocol.VerifierResult, error) { + return execute(ctx, r, r.verificationsPolicies, func() (map[protocol.Bytes32]protocol.VerifierResult, error) { return r.underlying.GetVerifications(ctx, messageIDs) }) } -func execute[T any](r *ResilientReader, policies executorPolicies[T], fn func() (T, error)) (T, error) { - result, err := policies.executor.GetWithExecution(func(failsafe.Execution[T]) (T, error) { +func execute[T any](ctx context.Context, r *ResilientReader, policies executorPolicies[T], fn func() (T, error)) (T, error) { + result, err := policies.executor.WithContext(ctx).GetWithExecution(func(failsafe.Execution[T]) (T, error) { return fn() }) if err != nil { diff --git a/indexer/pkg/readers/resilient_reader_test.go b/indexer/pkg/readers/resilient_reader_test.go new file mode 100644 index 000000000..85e7cee3d --- /dev/null +++ b/indexer/pkg/readers/resilient_reader_test.go @@ -0,0 +1,130 @@ +package readers + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "github.com/smartcontractkit/chainlink-ccv/common" + "github.com/smartcontractkit/chainlink-ccv/indexer/pkg/config" + "github.com/smartcontractkit/chainlink-ccv/protocol" + "github.com/smartcontractkit/chainlink-common/pkg/logger" +) + +func TestNewResilienceConfig(t *testing.T) { + def := DefaultResilienceConfig() + + t.Run("zero config returns all defaults", func(t *testing.T) { + rc := NewResilienceConfig(config.ResilienceConfig{}) + assert.Equal(t, def, rc) + }) + + t.Run("partial overrides keep defaults for unset fields", func(t *testing.T) { + rc := NewResilienceConfig(config.ResilienceConfig{ + MaxRequestsPerSecond: 100, + RequestTimeout: common.Duration(30 * time.Second), + }) + assert.Equal(t, uint(100), rc.MaxRequestsPerSecond) + assert.Equal(t, 30*time.Second, rc.RequestTimeout) + assert.Equal(t, def.MaxConcurrentRequests, rc.MaxConcurrentRequests) + assert.Equal(t, def.FailureThreshold, rc.FailureThreshold) + assert.Equal(t, def.SuccessThreshold, rc.SuccessThreshold) + assert.Equal(t, def.CircuitBreakerDelay, rc.CircuitBreakerDelay) + assert.Equal(t, def.CircuitBreakerTimeout, rc.CircuitBreakerTimeout) + assert.Equal(t, def.MaxRetries, rc.MaxRetries) + assert.Equal(t, def.RetryDelay, rc.RetryDelay) + assert.Equal(t, def.RetryMaxDelay, rc.RetryMaxDelay) + }) + + t.Run("full overrides", func(t *testing.T) { + in := config.ResilienceConfig{ + MaxRequestsPerSecond: 50, + MaxConcurrentRequests: 20, + FailureThreshold: 10, + SuccessThreshold: 7, + CircuitBreakerDelay: common.Duration(5 * time.Second), + CircuitBreakerTimeout: common.Duration(2 * time.Second), + RequestTimeout: common.Duration(15 * time.Second), + MaxRetries: 5, + RetryDelay: common.Duration(500 * time.Millisecond), + RetryMaxDelay: common.Duration(30 * time.Second), + } + rc := NewResilienceConfig(in) + assert.Equal(t, uint(50), rc.MaxRequestsPerSecond) + assert.Equal(t, uint(20), rc.MaxConcurrentRequests) + assert.Equal(t, uint32(10), rc.FailureThreshold) + assert.Equal(t, uint32(7), rc.SuccessThreshold) + assert.Equal(t, 5*time.Second, rc.CircuitBreakerDelay) + assert.Equal(t, 2*time.Second, rc.CircuitBreakerTimeout) + assert.Equal(t, 15*time.Second, rc.RequestTimeout) + assert.Equal(t, 5, rc.MaxRetries) + assert.Equal(t, 500*time.Millisecond, rc.RetryDelay) + assert.Equal(t, 30*time.Second, rc.RetryMaxDelay) + }) +} + +type mockOffchainReader struct { + mu sync.Mutex + callCount int + responses []protocol.QueryResponse +} + +func (m *mockOffchainReader) ReadCCVData(ctx context.Context) ([]protocol.QueryResponse, error) { + m.mu.Lock() + m.callCount++ + m.mu.Unlock() + return m.responses, nil +} + +func (m *mockOffchainReader) GetVerifications(ctx context.Context, messageIDs []protocol.Bytes32) (map[protocol.Bytes32]protocol.VerifierResult, error) { + return nil, nil +} + +func (m *mockOffchainReader) getCallCount() int { + m.mu.Lock() + defer m.mu.Unlock() + return m.callCount +} + +func TestResilientReader_RetryOnRateLimit(t *testing.T) { + mock := &mockOffchainReader{ + responses: []protocol.QueryResponse{{}}, + } + lggr, err := logger.New() + require.NoError(t, err) + + cfg := ResilienceConfig{ + FailureThreshold: 100, + SuccessThreshold: 3, + CircuitBreakerDelay: 3 * time.Second, + CircuitBreakerTimeout: 1 * time.Second, + RequestTimeout: 10 * time.Second, + MaxConcurrentRequests: 5, + MaxRequestsPerSecond: 1, + MaxRetries: 5, + RetryDelay: 50 * time.Millisecond, + RetryMaxDelay: 500 * time.Millisecond, + } + + rr := NewResilientReader(mock, lggr, cfg) + ctx := context.Background() + + resp1, err := rr.ReadCCVData(ctx) + require.NoError(t, err) + assert.Len(t, resp1, 1) + + start := time.Now() + resp2, err := rr.ReadCCVData(ctx) + elapsed := time.Since(start) + require.NoError(t, err) + assert.Len(t, resp2, 1) + + assert.Greater(t, elapsed, 500*time.Millisecond, + "second call should have waited for rate limit window reset via retries") + assert.Equal(t, 2, mock.getCallCount(), + "mock should be called twice — rate-limited attempts don't reach downstream") +} diff --git a/indexer/pkg/readers/rest_reader.go b/indexer/pkg/readers/rest_reader.go index f726acf74..02a22f115 100644 --- a/indexer/pkg/readers/rest_reader.go +++ b/indexer/pkg/readers/rest_reader.go @@ -37,6 +37,7 @@ type RestReaderConfig struct { HTTPClient *http.Client // Custom HTTP client (optional) Logger logger.Logger // Logger instance (required) Metrics common.IndexerMetricLabeler // Metrics instance (required) + Resilience ResilienceConfig // Resilience policies (zero values use defaults) } type restReader struct { @@ -68,7 +69,7 @@ func NewRestReader(config RestReaderConfig) *ResilientReader { m: config.Metrics, } - return NewResilientReader(NewObservedReader(underlying, nil, config.Metrics), config.Logger, DefaultResilienceConfig()) + return NewResilientReader(NewObservedReader(underlying, nil, config.Metrics), config.Logger, config.Resilience) } func (r *restReader) GetVerifications(ctx context.Context, messageIDs []protocol.Bytes32) (map[protocol.Bytes32]protocol.VerifierResult, error) { diff --git a/indexer/pkg/readers/rest_reader_test.go b/indexer/pkg/readers/rest_reader_test.go index 8dc23bc7c..2463b6ccc 100644 --- a/indexer/pkg/readers/rest_reader_test.go +++ b/indexer/pkg/readers/rest_reader_test.go @@ -116,6 +116,7 @@ func TestRestReader_GetVerifications_404_ReturnsEmptyMapAndNoError(t *testing.T) HTTPClient: server.Client(), Logger: lggr, Metrics: monitoring.NewNoopIndexerMetricLabeler(), + Resilience: DefaultResilienceConfig(), }) result, err := rr.GetVerifications(context.Background(), []protocol.Bytes32{messageID}) @@ -140,6 +141,7 @@ func TestRestReader_GetVerifications_404_MalformedBody_ReturnsEmptyMapAndNoError HTTPClient: server.Client(), Logger: lggr, Metrics: monitoring.NewNoopIndexerMetricLabeler(), + Resilience: DefaultResilienceConfig(), }) messageID := protocol.Bytes32{1, 2, 3} @@ -169,6 +171,7 @@ func TestRestReader_GetVerifications_404_DoesNotOpenCircuitBreaker(t *testing.T) HTTPClient: server.Client(), Logger: lggr, Metrics: monitoring.NewNoopIndexerMetricLabeler(), + Resilience: DefaultResilienceConfig(), }) ctx := context.Background()