mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-27 23:26:44 +00:00
consensus/bor/heimdall: add new meter and timer metrics in heimdall client (#459)
* heimdall client: add new meter and timer metrics * fix: error in handling valid and invalid request metric * use ctx outside of request * add separate metrics for each request * fix: modify metric name * fix: renmae metric name by removing hyphen * use ctx value for request type * fix linters * fix panic on using string as key for ctx and on nil return * refactor * refactor * rm comment Co-authored-by: Evgeny Danienko <6655321@bk.ru>
This commit is contained in:
parent
2a92cb1ecc
commit
f574c4a3c6
2 changed files with 115 additions and 4 deletions
|
|
@ -15,6 +15,7 @@ import (
|
|||
"github.com/ethereum/go-ethereum/consensus/bor/heimdall/checkpoint"
|
||||
"github.com/ethereum/go-ethereum/consensus/bor/heimdall/span"
|
||||
"github.com/ethereum/go-ethereum/log"
|
||||
"github.com/ethereum/go-ethereum/metrics"
|
||||
)
|
||||
|
||||
var (
|
||||
|
|
@ -46,6 +47,12 @@ type HeimdallClient struct {
|
|||
closeCh chan struct{}
|
||||
}
|
||||
|
||||
type Request struct {
|
||||
client http.Client
|
||||
url *url.URL
|
||||
start time.Time
|
||||
}
|
||||
|
||||
func NewHeimdallClient(urlString string) *HeimdallClient {
|
||||
return &HeimdallClient{
|
||||
urlString: urlString,
|
||||
|
|
@ -76,6 +83,8 @@ func (h *HeimdallClient) StateSyncEvents(ctx context.Context, fromID uint64, to
|
|||
|
||||
log.Info("Fetching state sync events", "queryParams", url.RawQuery)
|
||||
|
||||
ctx = withRequestType(ctx, stateSyncRequest)
|
||||
|
||||
response, err := FetchWithRetry[StateSyncEventsResponse](ctx, h.client, url, h.closeCh)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -108,6 +117,8 @@ func (h *HeimdallClient) Span(ctx context.Context, spanID uint64) (*span.Heimdal
|
|||
return nil, err
|
||||
}
|
||||
|
||||
ctx = withRequestType(ctx, spanRequest)
|
||||
|
||||
response, err := FetchWithRetry[SpanResponse](ctx, h.client, url, h.closeCh)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -123,6 +134,8 @@ func (h *HeimdallClient) FetchCheckpoint(ctx context.Context, number int64) (*ch
|
|||
return nil, err
|
||||
}
|
||||
|
||||
ctx = withRequestType(ctx, checkpointRequest)
|
||||
|
||||
response, err := FetchWithRetry[checkpoint.CheckpointResponse](ctx, h.client, url, h.closeCh)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
|
|
@ -138,6 +151,8 @@ func (h *HeimdallClient) FetchCheckpointCount(ctx context.Context) (int64, error
|
|||
return 0, err
|
||||
}
|
||||
|
||||
ctx = withRequestType(ctx, checkpointCountRequest)
|
||||
|
||||
response, err := FetchWithRetry[checkpoint.CheckpointCountResponse](ctx, h.client, url, h.closeCh)
|
||||
if err != nil {
|
||||
return 0, err
|
||||
|
|
@ -149,7 +164,9 @@ func (h *HeimdallClient) FetchCheckpointCount(ctx context.Context) (int64, error
|
|||
// FetchWithRetry returns data from heimdall with retry
|
||||
func FetchWithRetry[T any](ctx context.Context, client http.Client, url *url.URL, closeCh chan struct{}) (*T, error) {
|
||||
// request data once
|
||||
result, err := Fetch[T](ctx, client, url)
|
||||
request := &Request{client: client, url: url, start: time.Now()}
|
||||
result, err := Fetch[T](ctx, request)
|
||||
|
||||
if err == nil {
|
||||
return result, nil
|
||||
}
|
||||
|
|
@ -181,7 +198,9 @@ retryLoop:
|
|||
|
||||
return nil, ErrShutdownDetected
|
||||
case <-ticker.C:
|
||||
result, err = Fetch[T](ctx, client, url)
|
||||
request = &Request{client: client, url: url, start: time.Now()}
|
||||
result, err = Fetch[T](ctx, request)
|
||||
|
||||
if err != nil {
|
||||
if attempt%logEach == 0 {
|
||||
log.Warn("an error while trying fetching from Heimdall", "attempt", attempt, "error", err)
|
||||
|
|
@ -196,10 +215,18 @@ retryLoop:
|
|||
}
|
||||
|
||||
// Fetch returns data from heimdall
|
||||
func Fetch[T any](ctx context.Context, client http.Client, url *url.URL) (*T, error) {
|
||||
func Fetch[T any](ctx context.Context, request *Request) (*T, error) {
|
||||
isSuccessful := false
|
||||
|
||||
defer func() {
|
||||
if metrics.EnabledExpensive {
|
||||
sendMetrics(ctx, request.start, isSuccessful)
|
||||
}
|
||||
}()
|
||||
|
||||
result := new(T)
|
||||
|
||||
body, err := internalFetchWithTimeout(ctx, client, url)
|
||||
body, err := internalFetchWithTimeout(ctx, request.client, request.url)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
|
@ -213,6 +240,8 @@ func Fetch[T any](ctx context.Context, client http.Client, url *url.URL) (*T, er
|
|||
return nil, err
|
||||
}
|
||||
|
||||
isSuccessful = true
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
|
|
|
|||
82
consensus/bor/heimdall/metrics.go
Normal file
82
consensus/bor/heimdall/metrics.go
Normal file
|
|
@ -0,0 +1,82 @@
|
|||
package heimdall
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/ethereum/go-ethereum/metrics"
|
||||
)
|
||||
|
||||
type (
|
||||
requestTypeKey struct{}
|
||||
requestType string
|
||||
|
||||
meter struct {
|
||||
request map[bool]metrics.Meter // map[isSuccessful]metrics.Meter
|
||||
timer metrics.Timer
|
||||
}
|
||||
)
|
||||
|
||||
const (
|
||||
stateSyncRequest requestType = "state-sync"
|
||||
spanRequest requestType = "span"
|
||||
checkpointRequest requestType = "checkpoint"
|
||||
checkpointCountRequest requestType = "checkpoint-count"
|
||||
)
|
||||
|
||||
func withRequestType(ctx context.Context, reqType requestType) context.Context {
|
||||
return context.WithValue(ctx, requestTypeKey{}, reqType)
|
||||
}
|
||||
|
||||
func getRequestType(ctx context.Context) (requestType, bool) {
|
||||
reqType, ok := ctx.Value(requestTypeKey{}).(requestType)
|
||||
return reqType, ok
|
||||
}
|
||||
|
||||
var (
|
||||
requestMeters = map[requestType]meter{
|
||||
stateSyncRequest: {
|
||||
request: map[bool]metrics.Meter{
|
||||
true: metrics.NewRegisteredMeter("client/requests/statesync/valid", nil),
|
||||
false: metrics.NewRegisteredMeter("client/requests/statesync/invalid", nil),
|
||||
},
|
||||
timer: metrics.NewRegisteredTimer("client/requests/statesync/duration", nil),
|
||||
},
|
||||
spanRequest: {
|
||||
request: map[bool]metrics.Meter{
|
||||
true: metrics.NewRegisteredMeter("client/requests/span/valid", nil),
|
||||
false: metrics.NewRegisteredMeter("client/requests/span/invalid", nil),
|
||||
},
|
||||
timer: metrics.NewRegisteredTimer("client/requests/span/duration", nil),
|
||||
},
|
||||
checkpointRequest: {
|
||||
request: map[bool]metrics.Meter{
|
||||
true: metrics.NewRegisteredMeter("client/requests/checkpoint/valid", nil),
|
||||
false: metrics.NewRegisteredMeter("client/requests/checkpoint/invalid", nil),
|
||||
},
|
||||
timer: metrics.NewRegisteredTimer("client/requests/checkpoint/duration", nil),
|
||||
},
|
||||
checkpointCountRequest: {
|
||||
request: map[bool]metrics.Meter{
|
||||
true: metrics.NewRegisteredMeter("client/requests/checkpointcount/valid", nil),
|
||||
false: metrics.NewRegisteredMeter("client/requests/checkpointcount/invalid", nil),
|
||||
},
|
||||
timer: metrics.NewRegisteredTimer("client/requests/checkpointcount/duration", nil),
|
||||
},
|
||||
}
|
||||
)
|
||||
|
||||
func sendMetrics(ctx context.Context, start time.Time, isSuccessful bool) {
|
||||
reqType, ok := getRequestType(ctx)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
meters, ok := requestMeters[reqType]
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
meters.request[isSuccessful].Mark(1)
|
||||
meters.timer.Update(time.Since(start))
|
||||
}
|
||||
Loading…
Reference in a new issue