Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
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
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
9 changes: 9 additions & 0 deletions indexer/config.example.toml
Original file line number Diff line number Diff line change
Expand Up @@ -52,3 +52,12 @@ LockTimeout = 30
[API]
[API.RateLimit]
Enabled = false

[Resilience]
MaxRequestsPerSecond = 5
MaxConcurrentRequests = 5
FailureThreshold = 5
SuccessThreshold = 3
CircuitBreakerDelay = "3s"
CircuitBreakerTimeout = "1s"
RequestTimeout = "10s"
28 changes: 28 additions & 0 deletions indexer/pkg/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,34 @@ type Config struct {
Storage StorageConfig `toml:"Storage"`
// API is the configuration for the API inside the indexer.
API APIConfig `toml:"API"`
// Resilience is the configuration for the resilient reader rate limiting.
Resilience ResilienceConfig `toml:"Resilience"`
Comment thread
huangzhen1997 marked this conversation as resolved.
Outdated
}

// 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"`
}

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
29 changes: 29 additions & 0 deletions indexer/pkg/readers/resilient_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import (
"github.com/failsafe-go/failsafe-go/ratelimiter"
"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 Down Expand Up @@ -52,6 +53,34 @@ func DefaultResilienceConfig() ResilienceConfig {
}
}

// 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)
}
return rc
}

type executorPolicies[T any] struct {
executor failsafe.Executor[T]
circuitBreaker circuitbreaker.CircuitBreaker[T]
Expand Down
54 changes: 54 additions & 0 deletions indexer/pkg/readers/resilient_reader_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
package readers

import (
"testing"
"time"

"github.com/stretchr/testify/assert"

"github.com/smartcontractkit/chainlink-ccv/common"
"github.com/smartcontractkit/chainlink-ccv/indexer/pkg/config"
)

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)
})

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),
}
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)
})
}
3 changes: 2 additions & 1 deletion indexer/pkg/readers/rest_reader.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down Expand Up @@ -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)
Comment thread
huangzhen1997 marked this conversation as resolved.
}

func (r *restReader) GetVerifications(ctx context.Context, messageIDs []protocol.Bytes32) (map[protocol.Bytes32]protocol.VerifierResult, error) {
Expand Down
3 changes: 3 additions & 0 deletions indexer/pkg/readers/rest_reader_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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})
Expand All @@ -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}
Expand Down Expand Up @@ -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()
Expand Down
Loading