From f574c4a3c6cc48b4b2102fe75c83403508c8a692 Mon Sep 17 00:00:00 2001 From: Manav Darji Date: Tue, 23 Aug 2022 18:00:54 +0530 Subject: [PATCH] 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> --- consensus/bor/heimdall/client.go | 37 ++++++++++++-- consensus/bor/heimdall/metrics.go | 82 +++++++++++++++++++++++++++++++ 2 files changed, 115 insertions(+), 4 deletions(-) create mode 100644 consensus/bor/heimdall/metrics.go diff --git a/consensus/bor/heimdall/client.go b/consensus/bor/heimdall/client.go index 0f17cbc4d4..b6daa7e8d8 100644 --- a/consensus/bor/heimdall/client.go +++ b/consensus/bor/heimdall/client.go @@ -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 } diff --git a/consensus/bor/heimdall/metrics.go b/consensus/bor/heimdall/metrics.go new file mode 100644 index 0000000000..99d7ca65ac --- /dev/null +++ b/consensus/bor/heimdall/metrics.go @@ -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)) +}