Skip to content
Merged
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
5 changes: 5 additions & 0 deletions cmd/evmconnect.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,11 @@ func run() error {
return err
}

// Emit the block listener metrics into the metrics registry of the manager
if err := c.BlockListener().InitMetrics(ctx, m.MetricsRegistry()); err != nil {
return err
}

// Setup signal handling to cancel the context, which shuts down the API Server
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
go func() {
Expand Down
1 change: 1 addition & 0 deletions internal/msgs/en_error_messages.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,4 +91,5 @@ var (
MsgTransactionEstimateTooLargeForBlock = ffe("FF23071", "Gas estimate %s (scaled at %.2f from estimate %s) too large for the current block gas limit %s")
MsgMonitoredHeadLengthInvalid = ffe("FF23072", "Monitored head length must be greater than or equal to 1 value=%d")
MsgUnknownJSONFormatOptionValue = ffe("FF23073", "Unknown value '%s' for JSON formatting option '%s'. Supported values: %s")
MsgMetricsInitFail = ffe("FF23074", "Failed to initialize metrics for subsystem '%s'")
)
19 changes: 19 additions & 0 deletions mocks/ethblocklistenermocks/block_listener.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

41 changes: 32 additions & 9 deletions pkg/ethblocklistener/blocklistener.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import (
"github.com/hyperledger-firefly/common/pkg/fftypes"
"github.com/hyperledger-firefly/common/pkg/i18n"
"github.com/hyperledger-firefly/common/pkg/log"
"github.com/hyperledger-firefly/common/pkg/metric"
"github.com/hyperledger-firefly/common/pkg/retry"
"github.com/hyperledger-firefly/common/pkg/wsclient"
"github.com/hyperledger-firefly/evmconnect/internal/msgs"
Expand Down Expand Up @@ -89,6 +90,7 @@ type BlockListener interface {
WaitClosed()
GetBackend() rpcbackend.RPC
UTSetBackend(rpcbackend.RPC)
InitMetrics(ctx context.Context, registry metric.MetricsRegistry) error
}

func toMinimalBlockInfoList(blocks []*ethrpc.BlockInfoJSONRPC) []*ethrpc.MinimalBlockInfo {
Expand Down Expand Up @@ -142,6 +144,10 @@ type blockListener struct {

// headBlockNumber mode: last head value sent on the block listener channel (only written from listenLoop)
currentChainHead uint64

// metrics are optional - only emitted once InitMetrics has been called
metricsLock sync.RWMutex
metrics metric.MetricsManager
}

func NewBlockListener(ctx context.Context, retry *retry.Retry, conf *BlockListenerConfig, httpConf *ffresty.Config, wsConf *wsclient.WSConfig) (bl BlockListener, err error) {
Expand Down Expand Up @@ -295,26 +301,28 @@ func (bl *blockListener) establishBlockHeightWithRetry() error {
}

// Now get the block height
var hexBlockHeight ethtypes.HexInteger
rpcErr := bl.backend.CallRPC(bl.ctx, &hexBlockHeight, "eth_blockNumber")
if rpcErr != nil {
log.L(bl.ctx).Warnf("Block height could not be obtained: %s", rpcErr.Message)
return true, rpcErr.Error()
head, err := bl.queryBlockHeightFromRPC()
if err != nil {
log.L(bl.ctx).Warnf("Block height could not be obtained: %s", err)
return true, err
}

bl.setHighestBlock(hexBlockHeight.BigInt().Uint64())
bl.setHighestBlock(head)
return false, nil
})
}

// refreshHighestBlockFromRPC updates highestBlock from eth_blockNumber. Caller must not hold canonicalChainLock.
func (bl *blockListener) refreshHighestBlockFromRPC() (uint64, error) {
// queryBlockHeightFromRPC queries eth_blockNumber and returns the result, without updating any listener
// state. Caller must not hold canonicalChainLock. The height the node reports is recorded on the target block height gauge.
func (bl *blockListener) queryBlockHeightFromRPC() (uint64, error) {
var hexBlockHeight ethtypes.HexInteger
rpcErr := bl.backend.CallRPC(bl.ctx, &hexBlockHeight, "eth_blockNumber")
if rpcErr != nil {
bl.incPollFailureMetric("eth_blockNumber")
return 0, rpcErr.Error()
}
head := hexBlockHeight.BigInt().Uint64()
bl.setBlockHeightMetric(metricTargetBlockHeight, head)
return head, nil
}

Expand Down Expand Up @@ -367,10 +375,17 @@ func (bl *blockListener) listenLoop() {
}
}

// In full chain tracking mode, the loop below never queries the height the node reports, so we refresh
// it here for the target metric. Done ahead of the filter calls.
if bl.ChainTrackingMode != ffcapi.ChainTrackingModeLight {
bl.refreshTargetBlockHeightMetric()
}

if filter == "" {
err := bl.backend.CallRPC(bl.ctx, &filter, "eth_newBlockFilter")
if err != nil {
log.L(bl.ctx).Errorf("Failed to establish new block filter: %s", err.Message)
bl.incPollFailureMetric("eth_newBlockFilter")
failCount++
continue
}
Expand All @@ -393,24 +408,28 @@ func (bl *blockListener) listenLoop() {
gapPotential = true
}
log.L(bl.ctx).Errorf("Failed to query block filter changes: %s", rpcErr.Message)
bl.incPollFailureMetric("eth_getFilterChanges")
failCount++
continue
}
log.L(bl.ctx).Debugf("Block filter received new block hashes: %+v", blockHashes)
}

if bl.ChainTrackingMode == ffcapi.ChainTrackingModeLight {
Comment thread
peterbroadhurst marked this conversation as resolved.
head, err := bl.refreshHighestBlockFromRPC()
head, err := bl.queryBlockHeightFromRPC()
if err != nil {
log.L(bl.ctx).Errorf("Failed to refresh chain head: %s", err)
failCount++
continue
}
// In light mode there is no canonical chain being built, so the head we dispatch to
// consumers is what we report as the canonical height
if head == bl.currentChainHead {
failCount = 0
continue
}
bl.currentChainHead = head
bl.setBlockHeightMetric(metricCanonicalBlockHeight, bl.currentChainHead)
update := &ffcapi.BlockHashEvent{GapPotential: false, Created: fftypes.Now(), HeadBlockNumber: bl.currentChainHead}
bl.consumerMux.Lock()
consumers := make([]*BlockUpdateConsumer, 0, len(bl.consumers))
Expand Down Expand Up @@ -778,6 +797,7 @@ func (bl *blockListener) GetHeadBlockNumber(_ context.Context) uint64 {
}

func (bl *blockListener) setHighestBlock(block uint64) {
defer bl.setBlockHeightMetric(metricCanonicalBlockHeight, block)
bl.canonicalChainLock.Lock()
defer bl.canonicalChainLock.Unlock()
bl.highestBlock = block
Expand All @@ -793,6 +813,9 @@ func (bl *blockListener) checkAndSetHighestBlock(bi *ethrpc.BlockInfoJSONRPC) {
bl.highestBlock = block
bl.highestBlockSet = true
bl.headBlockInfo = bi
// The gauge is bound to the same variable GetHighestBlock reports to event streams, so it is the
// head we are actually tracking rather than a separate sample of it.
bl.setBlockHeightMetric(metricCanonicalBlockHeight, block)
} else if block == bl.highestBlock {
// Height already known from eth_blockNumber. Store the first full block at that height.
bl.headBlockInfo = bi
Expand Down
92 changes: 92 additions & 0 deletions pkg/ethblocklistener/blocklistener_metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
// Copyright © 2026 Kaleido, Inc.
//
// SPDX-License-Identifier: Apache-2.0
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

package ethblocklistener

import (
"context"

"github.com/hyperledger-firefly/common/pkg/i18n"
"github.com/hyperledger-firefly/common/pkg/log"
"github.com/hyperledger-firefly/common/pkg/metric"
"github.com/hyperledger-firefly/evmconnect/internal/msgs"
)

const (
metricsSubsystem = "blocklistener"

// metricTargetBlockHeight is the block height the endpoint we are connected to reports via eth_blockNumber.
// Emitted from queryBlockHeightFromRPC, so it is always the value we last received from the node.
metricTargetBlockHeight = "target_block_height"
// metricCanonicalBlockHeight is the head of the chain this listener is tracking - in full chain tracking
// mode the head of the in-memory canonical chain built from the block filter / newHeads subscription,
// and in light mode the head we dispatch to consumers. It should track the target height very closely.
metricCanonicalBlockHeight = "canonical_block_height"
// metricPollFailures counts the JSON/RPC polls the listen loop makes that failed, labelled by method.
metricPollFailures = "poll_failures_total"
metricLabelPollFailures = "method"
)

// InitMetrics registers the block listener metrics against the supplied registry.
func (bl *blockListener) InitMetrics(ctx context.Context, registry metric.MetricsRegistry) error {
mm, err := registry.NewMetricsManagerForSubsystem(ctx, metricsSubsystem)
if err != nil {
return i18n.WrapError(ctx, err, msgs.MsgMetricsInitFail, metricsSubsystem)
}
mm.NewGaugeMetric(ctx, metricTargetBlockHeight, "The block height reported by the connected node via eth_blockNumber", false)
mm.NewGaugeMetric(ctx, metricCanonicalBlockHeight, "The block height of the head of the chain tracked by the block listener", false)
mm.NewCounterMetricWithLabels(ctx, metricPollFailures, "The number of block listener JSON/RPC polls that have failed, by method", []string{metricLabelPollFailures}, false)

bl.metricsLock.Lock()
defer bl.metricsLock.Unlock()

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The lock is needed because the lifecycle of a metrics registry is tied to a transaction manager, but the block listener is not initiated inside the transaction manager.

In order to avoid a lock, a standalone metrics registry needs to be created first and then passed into both the listener and the transaction manager. That approach will reflect the fact that transaction manager no longer controls the metric prefix for components that emit custom metrics into the metrics that are linked to its metrics endpoint.

bl.metrics = mm
return nil
}

func (bl *blockListener) getMetrics() metric.MetricsManager {
bl.metricsLock.RLock()
defer bl.metricsLock.RUnlock()
return bl.metrics
}

func (bl *blockListener) setBlockHeightMetric(metricName string, blockHeight uint64) {
mm := bl.getMetrics()
if mm == nil {
return
}
mm.SetGaugeMetric(bl.ctx, metricName, float64(blockHeight), nil)
}

func (bl *blockListener) incPollFailureMetric(method string) {
mm := bl.getMetrics()
if mm == nil {
return
}
mm.IncCounterMetricWithLabels(bl.ctx, metricPollFailures, map[string]string{metricLabelPollFailures: method}, nil)
}

// refreshTargetBlockHeightMetric queries the node for the height it reports, purely so the target gauge
// stays current. Only needed in full chain tracking mode.
func (bl *blockListener) refreshTargetBlockHeightMetric() {
if bl.getMetrics() == nil {
return // never drive any query of the node when metrics are not enabled

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
return // never drive any query of the node when metrics are not enabled
return // never drive any query of the node when the metrics registry is not set

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This highlights a problem I overlooked as part of hyperledger-firefly/transaction-manager#160

Since the raw MetricsRegistry interface from ff-common was exposed, it didn't include information about whether the metrics are enabled in the config.

The control was wrapped in the internal Metrics interface: https://github.com/hyperledger-firefly/transaction-manager/blob/162a945b953601c8ea0811dc701a4d4e06d348da/internal/metrics/metrics.go#L90.

But feels like a separate PR to sort this out, as there is a feature gap in the metrics registry

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

}
if _, err := bl.queryBlockHeightFromRPC(); err != nil {
// Diagnostic only - the failure is recorded on the query failure counter, and the listen loop
// has its own error handling for the chain state
log.L(bl.ctx).Warnf("Failed to refresh target block height: %s", err)
}
}
Loading