From abc5242804f7b8dd70ec782066f2f7767a26496e Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Wed, 14 Dec 2022 03:52:03 +0300 Subject: [PATCH 1/3] fix daemon heads + head metrics --- .../dshackle/upstream/AbstractHead.kt | 68 +++++++++++++++---- .../ethereum_pos/EthereumPosMultiStream.kt | 42 ++++++------ 2 files changed, 74 insertions(+), 36 deletions(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt index 2e6aa454..cc1aa397 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt @@ -28,20 +28,26 @@ import reactor.core.publisher.Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER import reactor.core.publisher.Sinks.EmitResult.OK import reactor.core.scheduler.Schedulers import reactor.kotlin.core.publisher.toMono +import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.Executors +import java.util.concurrent.Future import java.util.concurrent.TimeUnit import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicInteger import java.util.concurrent.locks.ReentrantLock abstract class AbstractHead @JvmOverloads constructor( private val forkChoice: ForkChoice, private val blockValidator: BlockValidator = BlockValidator.ALWAYS_VALID, - awaitHeadTimeoutMs: Long = 60_000, + private val awaitHeadTimeoutMs: Long = 60_000, private val upstreamId: String = "" ) : Head { companion object { private val log = LoggerFactory.getLogger(AbstractHead::class.java) + private val instances: MutableMap = ConcurrentHashMap() + private val running: MutableMap = ConcurrentHashMap() + private val executor = Executors.newSingleThreadScheduledExecutor() } private var stream = Sinks.many().multicast().directBestEffort() @@ -50,25 +56,28 @@ abstract class AbstractHead @JvmOverloads constructor( private var stopping = false private var lastHeadUpdated = 0L private val lock = ReentrantLock() + private var future: Future<*>? = null + private val delayed = AtomicBoolean(false) init { - val state = AtomicBoolean(false) - Gauge.builder("stuck_head", state) { + + val className = this.javaClass.simpleName + + Gauge.builder("stuck_head", delayed) { if (it.get()) 1.0 else 0.0 - }.tag("upstream", upstreamId).tag("class", this.javaClass.simpleName).register(Metrics.globalRegistry) + }.tag("upstream", upstreamId).tag("class", className).register(Metrics.globalRegistry) Gauge.builder("current_head", forkChoice) { it.getHead()?.height?.toDouble() ?: 0.0 - }.tag("upstream", upstreamId).tag("class", this.javaClass.simpleName).register(Metrics.globalRegistry) - Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate( - { - val delay = System.currentTimeMillis() - lastHeadUpdated - val delayed = delay > awaitHeadTimeoutMs - state.set(delayed) - if (delayed) { - log.warn("No head updates $upstreamId for $delay ms @ ${this.javaClass}") - } - }, 300, 30, TimeUnit.SECONDS - ) + }.tag("upstream", upstreamId).tag("class", className).register(Metrics.globalRegistry) + + instances.computeIfAbsent(className) { + AtomicInteger(0).also { toHeadCountMetric(it, "allocated") } + }.incrementAndGet() + } + + protected fun finalize() { + instances[this.javaClass.simpleName]?.decrementAndGet() + println("Finish") } fun follow(source: Flux): Disposable { @@ -154,9 +163,38 @@ abstract class AbstractHead @JvmOverloads constructor( override fun stop() { stopping = true + log.debug("Stop ${this.javaClass.simpleName} $upstreamId") + future?.let { + running[this.javaClass.simpleName]?.decrementAndGet() + it.cancel(true) + } + future = null } override fun start() { stopping = false + log.debug("Start ${this.javaClass.simpleName} $upstreamId") + if (future == null) { + future = executor.scheduleAtFixedRate( + { + val delay = System.currentTimeMillis() - lastHeadUpdated + delayed.set(delay > awaitHeadTimeoutMs) + if (delayed.get()) { + log.warn("No head updates $upstreamId for $delay ms @ ${this.javaClass.simpleName}") + } + System.gc() + }, 180, 30, TimeUnit.SECONDS + ) + + running.computeIfAbsent(this.javaClass.simpleName) { + AtomicInteger(0).also { toHeadCountMetric(it, "running") } + }.incrementAndGet() + } + } + + private fun toHeadCountMetric(counter: AtomicInteger, status: String) { + Gauge.builder("head_count", counter) { + it.get().toDouble() + }.tag("class", this.javaClass.simpleName).tag("status", status).register(Metrics.globalRegistry) } } 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 abf12c63..5070dcb0 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 @@ -112,33 +112,33 @@ open class EthereumPosMultiStream( } override fun updateHead(): Head { - head?.let { - if (it is Lifecycle) { - it.stop() - } - } + this.head?.takeIf { it is Lifecycle }?.apply { stop() } lagObserver?.stop() lagObserver = null - val head = if (upstreams.size == 1) { - val upstream = upstreams.first() - upstream.setLag(0) - upstream.getHead().apply { - if (this is Lifecycle) { - this.start() + + return when (upstreams.size) { + 0 -> EmptyHead() + 1 -> upstreams.first().let { + it.setLag(0) + it.getHead().apply { + if (this is Lifecycle) { + start() + } } } - } else { - val heads = upstreams.map { it.getHead() } - val newHead = MergedHead(heads, PriorityForkChoice(), "ETH Pos Multistream").apply { - this.start() + + else -> upstreams.map { it.getHead() }.let { heads -> + MergedHead(heads, PriorityForkChoice(), "ETH Pos Multistream").apply { + start() + }.also { + this.lagObserver = EthereumPosHeadLagObserver(it, ArrayList(upstreams)).apply { + start() + } + } } - val lagObserver = EthereumPosHeadLagObserver(newHead, ArrayList(upstreams)) - this.lagObserver = lagObserver - lagObserver.start() - newHead + }.apply { + onHeadUpdated(this) } - onHeadUpdated(head) - return head } override fun getLabels(): Collection { From 833216b5d383c8cbd6374d756a4a136c6d5f9842 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Wed, 14 Dec 2022 12:59:47 +0300 Subject: [PATCH 2/3] remove System.gc --- src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt | 1 - 1 file changed, 1 deletion(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt index cc1aa397..3fa980a7 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt @@ -182,7 +182,6 @@ abstract class AbstractHead @JvmOverloads constructor( if (delayed.get()) { log.warn("No head updates $upstreamId for $delay ms @ ${this.javaClass.simpleName}") } - System.gc() }, 180, 30, TimeUnit.SECONDS ) From 0742c5f542ac08792f25c61193373435734a5303 Mon Sep 17 00:00:00 2001 From: Maksim Fomenkov Date: Wed, 14 Dec 2022 13:02:48 +0300 Subject: [PATCH 3/3] remove log --- src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt | 1 - 1 file changed, 1 deletion(-) diff --git a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt index 3fa980a7..216c7907 100644 --- a/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt +++ b/src/main/kotlin/io/emeraldpay/dshackle/upstream/AbstractHead.kt @@ -77,7 +77,6 @@ abstract class AbstractHead @JvmOverloads constructor( protected fun finalize() { instances[this.javaClass.simpleName]?.decrementAndGet() - println("Finish") } fun follow(source: Flux): Disposable {