From afa374547872e03b3268a3222deb62dddc5c4770 Mon Sep 17 00:00:00 2001 From: Termina1 Date: Sat, 3 Dec 2022 15:52:54 +0200 Subject: [PATCH 1/2] more logs for heads and better understanding of what they belong to --- .../io/emeraldpay/dshackle/upstream/AbstractHead.kt | 10 +++++----- .../io/emeraldpay/dshackle/upstream/MergedHead.kt | 11 +++++++---- .../dshackle/upstream/ethereum/EthereumMultistream.kt | 11 +++++------ .../ethereum/connectors/EthereumRpcConnector.kt | 2 +- .../upstream/ethereum_pos/EthereumPosMultiStream.kt | 8 ++++---- 5 files changed, 22 insertions(+), 20 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt index 0c90ff18..6ee6c0b4 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt @@ -53,7 +53,7 @@ abstract class AbstractHead @JvmOverloads constructor( { val delay = System.currentTimeMillis() - lastHeadUpdated if (delay > awaitHeadTimeoutMs) { - log.warn("No head updates for $delay ms @ ${this.javaClass} - restart") + log.warn("No head updates $upstreamId for $delay ms @ ${this.javaClass} - restart") if (lock.tryLock()) { try { start() @@ -82,10 +82,10 @@ abstract class AbstractHead @JvmOverloads constructor( // but technically it should never happen during normal work, only when the Head // is stopping if (it == SignalType.ON_ERROR && !stopping) { - log.warn("Received signal $it unexpectedly - restart head") + log.warn("Received signal $upstreamId $it unexpectedly - restart head") lastHeadUpdated = 0L } else { - log.warn("Received signal $it - stop emit new head!!!") + log.warn("Received signal $upstreamId $it - stop emit new head!!!") completed = true stream.tryEmitComplete() } @@ -95,7 +95,7 @@ abstract class AbstractHead @JvmOverloads constructor( val valid = runCatching { blockValidator.isValid(forkChoice.getHead(), block) }.onFailure { - log.error("Block ${block.hash} validation failed with '${it.message}'", it) + log.error("Block $upstreamId ${block.hash} validation failed with '${it.message}'", it) }.getOrElse { false } if (valid) { notifyBeforeBlock() @@ -113,7 +113,7 @@ abstract class AbstractHead @JvmOverloads constructor( is ForkChoice.ChoiceResult.Same -> {} } } else { - log.warn("Invalid block $block}") + log.warn("Invalid block $upstreamId $block}") } } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt index 7ea7d7f8..18e62a55 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/MergedHead.kt @@ -24,10 +24,11 @@ import org.slf4j.LoggerFactory import reactor.core.Disposable import reactor.core.publisher.Flux -class MergedHead( +class MergedHead @JvmOverloads constructor( private val sources: Iterable, - forkChoice: ForkChoice -) : AbstractHead(forkChoice), Lifecycle, CachesEnabled { + forkChoice: ForkChoice, + private val label: String = "" +) : AbstractHead(forkChoice, upstreamId = label), Lifecycle, CachesEnabled { companion object { private val log = LoggerFactory.getLogger(MergedHead::class.java) @@ -48,7 +49,9 @@ class MergedHead( } subscription?.dispose() subscription = super.follow( - Flux.merge(sources.map { it.getFlux() }).doOnNext { log.debug("New MERGED head $it") } + Flux.merge(sources.map { it.getFlux() }).doOnNext { + log.debug("New MERGED $label head $it") + } ) } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt index 69bb9ac7..899d5358 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/EthereumMultistream.kt @@ -128,7 +128,7 @@ open class EthereumMultistream( } } else { val heads = upstreams.map { it.getHead() } - val newHead = MergedHead(heads, MostWorkForkChoice()).apply { + val newHead = MergedHead(heads, MostWorkForkChoice(), "ETH Multistream").apply { this.start() } val lagObserver = EthereumHeadLagObserver(newHead, upstreams as Collection) @@ -166,13 +166,12 @@ open class EthereumMultistream( upstreams.filter { mather.matches(it) } .apply { log.debug("Found $size upstreams matching [${mather.describeInternal()}]") - } - .map { it.getHead() } - .let { + }.let { + val selected = it.map { it.getHead() } when (it.size) { 0 -> EmptyHead() - 1 -> it.first() - else -> MergedHead(it, MostWorkForkChoice()).apply { + 1 -> selected.first() + else -> MergedHead(selected, MostWorkForkChoice(), "Eth head ${it.map { it.getId() }}").apply { start() } } diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt index 21fefba1..4f61b43d 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum/connectors/EthereumRpcConnector.kt @@ -38,7 +38,7 @@ class EthereumRpcConnector( val wsHead = EthereumWsHead(conn, id, forkChoice, blockValidator) // receive bew blocks through WebSockets, but also periodically verify with RPC in case if WS failed val rpcHead = EthereumRpcHead(directReader, forkChoice, id, blockValidator, Duration.ofSeconds(60)) - head = MergedHead(listOf(rpcHead, wsHead), forkChoice) + head = MergedHead(listOf(rpcHead, wsHead), forkChoice, "Merged for $id") } else { conn = null log.warn("Setting up connector for $id upstream with RPC-only access, less effective than WS+RPC") diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt index 252ce212..bc4364fc 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/ethereum_pos/EthereumPosMultiStream.kt @@ -123,7 +123,7 @@ open class EthereumPosMultiStream( } } else { val heads = upstreams.map { it.getHead() } - val newHead = MergedHead(heads, PriorityForkChoice()).apply { + val newHead = MergedHead(heads, PriorityForkChoice(), "ETH Pos Multistream").apply { this.start() } val lagObserver = EthereumPosHeadLagObserver(newHead, upstreams as Collection) @@ -161,12 +161,12 @@ open class EthereumPosMultiStream( .apply { log.debug("Found $size upstreams matching [${mather.describeInternal()}]") } - .map { it.getHead() } .let { + val selected = it.map { it.getHead() } when (it.size) { 0 -> EmptyHead() - 1 -> it.first() - else -> MergedHead(it, PriorityForkChoice()).apply { + 1 -> selected.first() + else -> MergedHead(selected, PriorityForkChoice(), "ETH head for ${it.map { it.getId() }}").apply { start() } } From 5e15096da50af43016c1df8860112b493ece598f Mon Sep 17 00:00:00 2001 From: Termina1 Date: Sat, 3 Dec 2022 15:54:05 +0200 Subject: [PATCH 2/2] fix priority forkchoice, in case some nodes are faster than others --- .../kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt index 6ee6c0b4..7c7bb5b9 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt @@ -73,10 +73,10 @@ abstract class AbstractHead @JvmOverloads constructor( completed = false } return source - .distinctUntilChanged { - it.hash + .filter { + log.debug("Filtering block $upstreamId block $it") + forkChoice.filter(it) } - .filter { forkChoice.filter(it) } .doFinally { // close internal stream if upstream is finished, otherwise it gets stuck, // but technically it should never happen during normal work, only when the Head