Skip to content
Open
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
34 changes: 34 additions & 0 deletions docs/config/indexer/config.documented.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"

13 changes: 7 additions & 6 deletions indexer/cmd/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}
Expand All @@ -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
}
Expand All @@ -182,20 +182,21 @@ 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,
RequestTimeout: time.Duration(cfg.RequestTimeout),
MaxResponseBytes: config.EffectiveMaxResponseBytes(cfg.MaxResponseBytes),
Logger: lggr,
Metrics: m,
Resilience: readers.NewResilienceConfig(resiConfig),
}), nil
default:
return nil, errors.New("unknown verifier type")
Expand Down Expand Up @@ -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
Expand Down
9 changes: 5 additions & 4 deletions indexer/cmd/replay/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand All @@ -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))
}
}

Expand Down Expand Up @@ -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

Expand All @@ -428,14 +428,15 @@ 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,
RequestTimeout: time.Duration(vc.RequestTimeout),
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))
Expand Down
12 changes: 12 additions & 0 deletions indexer/config.example.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
38 changes: 38 additions & 0 deletions indexer/pkg/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"`
Comment on lines +77 to +79
// 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"`
Comment on lines +86 to +91
}

type SchedulerConfig struct {
Expand Down
9 changes: 4 additions & 5 deletions indexer/pkg/readers/aggregator_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Comment thread
huangzhen1997 marked this conversation as resolved.
}

// aggregatorRetryPolicyErrorHandler determines if an error from the aggregator should be retried.
Expand Down
89 changes: 76 additions & 13 deletions indexer/pkg/readers/resilient_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand All @@ -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.
Expand All @@ -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]
Expand All @@ -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)
}).
Expand All @@ -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)
Expand All @@ -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 {
Expand Down
Loading
Loading