|
|
|
@ -4,6 +4,7 @@ import ( |
|
|
|
|
"context" |
|
|
|
|
"fmt" |
|
|
|
|
"math/big" |
|
|
|
|
"reflect" |
|
|
|
|
"time" |
|
|
|
|
|
|
|
|
|
"encoding/hex" |
|
|
|
@ -42,6 +43,10 @@ type PublicBlockchainService struct { |
|
|
|
|
limiter *rate.Limiter |
|
|
|
|
rpcBlockFactory rpc_common.BlockFactory |
|
|
|
|
helper *bcServiceHelper |
|
|
|
|
// TEMP SOLUTION to rpc node spamming issue
|
|
|
|
|
limiterGetStakingNetworkInfo *rate.Limiter |
|
|
|
|
limiterGetSuperCommittees *rate.Limiter |
|
|
|
|
limiterGetCurrentUtilityMetrics *rate.Limiter |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
const ( |
|
|
|
@ -53,7 +58,7 @@ const ( |
|
|
|
|
func NewPublicBlockchainAPI(hmy *hmy.Harmony, version Version, limiterEnable bool, limit int) rpc.API { |
|
|
|
|
var limiter *rate.Limiter |
|
|
|
|
if limiterEnable { |
|
|
|
|
limiter := rate.NewLimiter(rate.Limit(limit), 1) |
|
|
|
|
limiter = rate.NewLimiter(rate.Limit(limit), limit) |
|
|
|
|
strLimit := fmt.Sprintf("%d", int64(limiter.Limit())) |
|
|
|
|
rpcRateLimitCounterVec.With(prometheus.Labels{ |
|
|
|
|
"rate_limit": strLimit, |
|
|
|
@ -61,9 +66,12 @@ func NewPublicBlockchainAPI(hmy *hmy.Harmony, version Version, limiterEnable boo |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
s := &PublicBlockchainService{ |
|
|
|
|
hmy: hmy, |
|
|
|
|
version: version, |
|
|
|
|
limiter: limiter, |
|
|
|
|
hmy: hmy, |
|
|
|
|
version: version, |
|
|
|
|
limiter: limiter, |
|
|
|
|
limiterGetStakingNetworkInfo: rate.NewLimiter(5, 10), |
|
|
|
|
limiterGetSuperCommittees: rate.NewLimiter(5, 10), |
|
|
|
|
limiterGetCurrentUtilityMetrics: rate.NewLimiter(5, 10), |
|
|
|
|
} |
|
|
|
|
s.helper = s.newHelper() |
|
|
|
|
|
|
|
|
@ -144,18 +152,19 @@ func (s *PublicBlockchainService) BlockNumber(ctx context.Context) (interface{}, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func (s *PublicBlockchainService) wait(ctx context.Context) error { |
|
|
|
|
if s.limiter != nil { |
|
|
|
|
func (s *PublicBlockchainService) wait(limiter *rate.Limiter, ctx context.Context) error { |
|
|
|
|
if limiter != nil { |
|
|
|
|
deadlineCtx, cancel := context.WithTimeout(ctx, DefaultRateLimiterWaitTimeout) |
|
|
|
|
defer cancel() |
|
|
|
|
if !s.limiter.Allow() { |
|
|
|
|
strLimit := fmt.Sprintf("%d", int64(s.limiter.Limit())) |
|
|
|
|
if !limiter.Allow() { |
|
|
|
|
strLimit := fmt.Sprintf("%d", int64(limiter.Limit())) |
|
|
|
|
name := reflect.TypeOf(limiter).Elem().Name() |
|
|
|
|
rpcRateLimitCounterVec.With(prometheus.Labels{ |
|
|
|
|
"rate_limit": strLimit, |
|
|
|
|
name: strLimit, |
|
|
|
|
}).Inc() |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
return s.limiter.Wait(deadlineCtx) |
|
|
|
|
return limiter.Wait(deadlineCtx) |
|
|
|
|
} |
|
|
|
|
return nil |
|
|
|
|
} |
|
|
|
@ -169,7 +178,7 @@ func (s *PublicBlockchainService) GetBlockByNumber( |
|
|
|
|
timer := DoMetricRPCRequest(GetBlockByNumber) |
|
|
|
|
defer DoRPCRequestDuration(GetBlockByNumber, timer) |
|
|
|
|
|
|
|
|
|
err = s.wait(ctx) |
|
|
|
|
err = s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetBlockByNumber, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -222,7 +231,7 @@ func (s *PublicBlockchainService) GetBlockByHash( |
|
|
|
|
timer := DoMetricRPCRequest(GetBlockByHash) |
|
|
|
|
defer DoRPCRequestDuration(GetBlockByHash, timer) |
|
|
|
|
|
|
|
|
|
err = s.wait(ctx) |
|
|
|
|
err = s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetBlockByHash, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -512,7 +521,7 @@ func (s *PublicBlockchainService) GetShardingStructure( |
|
|
|
|
timer := DoMetricRPCRequest(GetShardingStructure) |
|
|
|
|
defer DoRPCRequestDuration(GetShardingStructure, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetShardingStructure, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -569,7 +578,7 @@ func (s *PublicBlockchainService) LatestHeader(ctx context.Context) (StructuredR |
|
|
|
|
timer := DoMetricRPCRequest(LatestHeader) |
|
|
|
|
defer DoRPCRequestDuration(LatestHeader, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(LatestHeader, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -604,7 +613,7 @@ func (s *PublicBlockchainService) GetLastCrossLinks( |
|
|
|
|
timer := DoMetricRPCRequest(GetLastCrossLinks) |
|
|
|
|
defer DoRPCRequestDuration(GetLastCrossLinks, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetLastCrossLinks, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -642,7 +651,7 @@ func (s *PublicBlockchainService) GetHeaderByNumber( |
|
|
|
|
timer := DoMetricRPCRequest(GetHeaderByNumber) |
|
|
|
|
defer DoRPCRequestDuration(GetHeaderByNumber, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetHeaderByNumber, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -813,7 +822,7 @@ func (s *PublicBlockchainService) GetCurrentUtilityMetrics( |
|
|
|
|
timer := DoMetricRPCRequest(GetCurrentUtilityMetrics) |
|
|
|
|
defer DoRPCRequestDuration(GetCurrentUtilityMetrics, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiterGetCurrentUtilityMetrics, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetCurrentUtilityMetrics, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -842,7 +851,7 @@ func (s *PublicBlockchainService) GetSuperCommittees( |
|
|
|
|
timer := DoMetricRPCRequest(GetSuperCommittees) |
|
|
|
|
defer DoRPCRequestDuration(GetSuperCommittees, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiterGetSuperCommittees, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetSuperCommittees, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -871,7 +880,7 @@ func (s *PublicBlockchainService) GetCurrentBadBlocks( |
|
|
|
|
timer := DoMetricRPCRequest(GetCurrentBadBlocks) |
|
|
|
|
defer DoRPCRequestDuration(GetCurrentBadBlocks, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiter, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetCurrentBadBlocks, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
@ -913,7 +922,7 @@ func (s *PublicBlockchainService) GetStakingNetworkInfo( |
|
|
|
|
timer := DoMetricRPCRequest(GetStakingNetworkInfo) |
|
|
|
|
defer DoRPCRequestDuration(GetStakingNetworkInfo, timer) |
|
|
|
|
|
|
|
|
|
err := s.wait(ctx) |
|
|
|
|
err := s.wait(s.limiterGetStakingNetworkInfo, ctx) |
|
|
|
|
if err != nil { |
|
|
|
|
DoMetricRPCQueryInfo(GetStakingNetworkInfo, FailedNumber) |
|
|
|
|
return nil, err |
|
|
|
|