diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumFinalizationDetector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumFinalizationDetector.kt index 0c8722f6..2f9bf602 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumFinalizationDetector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumFinalizationDetector.kt @@ -112,4 +112,8 @@ class EthereumFinalizationDetector : FinalizationDetector { override fun getFinalizations(): Collection { return data.values } + + override fun getFinalizationByType(type: FinalizationType): FinalizationData? { + return data[type] + } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/finalization/FinalizationDetector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/finalization/FinalizationDetector.kt index 6563256b..8f540e44 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/finalization/FinalizationDetector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/finalization/FinalizationDetector.kt @@ -14,6 +14,8 @@ interface FinalizationDetector { fun getFinalizations(): Collection + fun getFinalizationByType(type: FinalizationType): FinalizationData? + fun addFinalization(finalization: FinalizationData) } @@ -30,6 +32,10 @@ class NoopFinalizationDetector : FinalizationDetector { return emptyList() } + override fun getFinalizationByType(type: FinalizationType): FinalizationData? { + return null + } + override fun addFinalization(finalization: FinalizationData) { } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/GenericUpstream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/GenericUpstream.kt index bd53fc3a..ea364ace 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/GenericUpstream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/generic/GenericUpstream.kt @@ -34,6 +34,9 @@ import io.emeraldpay.dshackle.upstream.generic.connectors.GenericConnector import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundData import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundServiceBuilder import io.emeraldpay.dshackle.upstream.lowerbound.LowerBoundType +import io.micrometer.core.instrument.Gauge +import io.micrometer.core.instrument.Meter +import io.micrometer.core.instrument.Metrics import org.springframework.context.Lifecycle import org.springframework.scheduling.concurrent.CustomizableThreadFactory import reactor.core.Disposable @@ -48,7 +51,7 @@ import java.util.function.Supplier open class GenericUpstream( id: String, - chain: Chain, + private val chain: Chain, hash: Short, options: ChainOptions.Options, role: UpstreamsConfig.UpstreamRole, @@ -106,6 +109,14 @@ open class GenericUpstream( private val hasLiveSubscriptionHead: AtomicBoolean = AtomicBoolean(getOptions().disableLivenessSubscriptionValidation) protected val connector: GenericConnector = connectorFactory.create(this, chain) + .also { upConnector -> + Gauge.builder("upstream_head", upConnector.getHead()) { + it.getCurrentHeight()?.toDouble() ?: 0.0 + } + .tag("upstreamId", id) + .tag("chain", chain.chainCode) + .register(Metrics.globalRegistry) + } private val livenessSubscription = AtomicReference() private val settingsDetector = upstreamSettingsDetectorBuilder(chain, this) private var rpcMethodsDetector: UpstreamRpcMethodsDetector? = null @@ -121,6 +132,8 @@ open class GenericUpstream( private val headLivenessState = Sinks.many().multicast().directBestEffort() + private val metrics = mutableMapOf() + override fun getHead(): Head { return connector.getHead() } @@ -387,7 +400,18 @@ open class GenericUpstream( finalizationDetectorSubscription.set( finalizationDetector.detectFinalization(this, chainConfig.expectedBlockTime, getChain()) .subscribeOn(finalizationScheduler) - .subscribe { + .subscribe { data -> + metrics.computeIfAbsent( + "${data.type}-${getId()}", + ) { + Gauge.builder("upstream_blocks", finalizationDetector) { detector -> + detector.getFinalizationByType(data.type)?.height?.toDouble() ?: 0.0 + } + .tag("blockType", data.type.toString()) + .tag("upstreamId", getId()) + .tag("chain", chain.chainCode) + .register(Metrics.globalRegistry) + } sendUpstreamStateEvent(UPDATED) }, )