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