better stuck head handling

This commit is contained in:
Maksim Fomenkov
2022-10-10 13:48:34 +03:00
parent 0aacc09862
commit c53913db6a
13 changed files with 90 additions and 42 deletions

View File

@@ -21,13 +21,17 @@ import org.slf4j.LoggerFactory
import reactor.core.Disposable
import reactor.core.publisher.Flux
import reactor.core.publisher.Sinks
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.atomic.AtomicLong
import java.util.concurrent.Executors
import java.util.concurrent.TimeUnit
abstract class AbstractHead(
private val forkChoice: ForkChoice,
private val blockValidator: BlockValidator = BlockValidator.ALWAYS_VALID
private val blockValidator: BlockValidator = BlockValidator.ALWAYS_VALID,
awaitHeadTimeoutMs: Long = 60_000
) : Head {
companion object {
@@ -37,7 +41,20 @@ abstract class AbstractHead(
private var stream = Sinks.many().multicast().directBestEffort<BlockContainer>()
private var completed = false
private val beforeBlockHandlers = ArrayList<Runnable>()
private val lastUpdateTime = AtomicLong(0L)
private var stopping = false
private var lastHeadUpdated = 0L
init {
Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(
{
val delay = System.currentTimeMillis() - lastHeadUpdated
if (delay > awaitHeadTimeoutMs) {
log.warn("No head updates for $delay ms @ ${this.javaClass} - restart")
start()
}
}, 300, 30, TimeUnit.SECONDS
)
}
fun follow(source: Flux<BlockContainer>): Disposable {
if (completed) {
@@ -54,8 +71,14 @@ abstract class AbstractHead(
// close internal stream if upstream is finished, otherwise it gets stuck,
// but technically it should never happen during normal work, only when the Head
// is stopping
completed = true
stream.tryEmitComplete()
if (stopping) {
log.info("Received signal $it - stop emit new head!!!")
completed = true
stream.tryEmitComplete()
} else {
log.warn("Received signal $it unexpectedly - restart head")
start()
}
}
.subscribeOn(Schedulers.boundedElastic())
.subscribe { block ->
@@ -69,12 +92,12 @@ abstract class AbstractHead(
when (val choiceResult = forkChoice.choose(block)) {
is ForkChoice.ChoiceResult.Updated -> {
val newHead = choiceResult.nwhead
log.debug("New block ${newHead.height} ${newHead.hash}")
val result = stream.tryEmitNext(newHead)
if (result.isFailure && result != Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER) {
log.warn("Failed to dispatch block: $result as ${this.javaClass}")
lastHeadUpdated = System.currentTimeMillis()
when (val result = stream.tryEmitNext(newHead)) {
OK -> log.debug("New block ${newHead.height} ${newHead.hash} @ ${this.javaClass}")
FAIL_ZERO_SUBSCRIBER -> log.debug("No subscribers for ${this.javaClass}")
else -> log.warn("Failed to dispatch block: $result as ${this.javaClass}")
}
lastUpdateTime.set(System.currentTimeMillis())
}
is ForkChoice.ChoiceResult.Same -> {}
@@ -100,7 +123,6 @@ abstract class AbstractHead(
}
override fun getFlux(): Flux<BlockContainer> {
val curHead = forkChoice.getHead()
return Flux.concat(
forkChoice.getHead().toMono(),
stream.asFlux()
@@ -115,6 +137,11 @@ abstract class AbstractHead(
return getCurrent()?.height
}
override fun getLastUpdateTime(): Long =
lastUpdateTime.get()
override fun stop() {
stopping = true
}
override fun start() {
stopping = false
}
}

View File

@@ -31,6 +31,9 @@ class EmptyHead : Head {
override fun getCurrentHeight(): Long? {
return null
}
override fun start() {
}
override fun getLastUpdateTime(): Long = 0L
override fun stop() {
}
}

View File

@@ -37,5 +37,8 @@ interface Head {
fun onBeforeBlock(handler: Runnable)
fun getCurrentHeight(): Long?
fun getLastUpdateTime(): Long
fun start()
fun stop()
}

View File

@@ -20,6 +20,7 @@ import com.google.common.annotations.VisibleForTesting
import io.emeraldpay.dshackle.cache.Caches
import io.emeraldpay.dshackle.cache.CachesEnabled
import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
import org.slf4j.LoggerFactory
import org.springframework.context.Lifecycle
import reactor.core.Disposable
import reactor.core.publisher.Flux
@@ -29,6 +30,10 @@ class MergedHead(
forkChoice: ForkChoice
) : AbstractHead(forkChoice), Lifecycle, CachesEnabled {
companion object {
private val log = LoggerFactory.getLogger(MergedHead::class.java)
}
private var subscription: Disposable? = null
override fun isRunning(): Boolean {
@@ -36,16 +41,20 @@ class MergedHead(
}
override fun start() {
super.start()
sources.forEach { head ->
if (head is Lifecycle && !head.isRunning) {
head.start()
}
}
subscription?.dispose()
subscription = super.follow(Flux.merge(sources.map { it.getFlux() }))
subscription = super.follow(
Flux.merge(sources.map { it.getFlux() }).doOnNext { log.debug("New MERGED head $it") }
)
}
override fun stop() {
super.stop()
sources.forEach { head ->
if (head is Lifecycle && head.isRunning) {
head.stop()

View File

@@ -36,7 +36,7 @@ class BitcoinRpcHead(
private val api: Reader<JsonRpcRequest, JsonRpcResponse>,
private val extractBlock: ExtractBlock,
private val interval: Duration = Duration.ofSeconds(15)
) : Head, AbstractHead(MostWorkForkChoice()), Lifecycle {
) : Head, AbstractHead(MostWorkForkChoice(), awaitHeadTimeoutMs = 1200_000), Lifecycle {
companion object {
private val log = LoggerFactory.getLogger(BitcoinRpcHead::class.java)
@@ -51,6 +51,7 @@ class BitcoinRpcHead(
}
override fun start() {
super.start()
if (refreshSubscription != null) {
log.warn("Called to start when running")
return
@@ -76,6 +77,7 @@ class BitcoinRpcHead(
}
override fun stop() {
super.stop()
val copy = refreshSubscription
refreshSubscription = null
copy?.dispose()

View File

@@ -21,7 +21,7 @@ class BitcoinZMQHead(
private val server: ZMQServer,
private val api: Reader<JsonRpcRequest, JsonRpcResponse>,
private val extractBlock: ExtractBlock,
) : Head, AbstractHead(MostWorkForkChoice()), Lifecycle {
) : Head, AbstractHead(MostWorkForkChoice(), awaitHeadTimeoutMs = 1200_000), Lifecycle {
companion object {
private val log = LoggerFactory.getLogger(BitcoinZMQHead::class.java)
@@ -51,11 +51,13 @@ class BitcoinZMQHead(
}
override fun start() {
super.start()
server.start()
refreshSubscription = super.follow(connect())
}
override fun stop() {
super.stop()
server.stop()
val copy = refreshSubscription
refreshSubscription = null

View File

@@ -48,6 +48,7 @@ class EthereumRpcHead(
private var refreshSubscription: Disposable? = null
override fun start() {
super.start()
refreshSubscription?.dispose()
val base = Flux.interval(interval)
.publishOn(scheduler)
@@ -62,6 +63,7 @@ class EthereumRpcHead(
}
override fun stop() {
super.stop()
refreshSubscription?.dispose()
refreshSubscription = null
}

View File

@@ -40,6 +40,7 @@ class EthereumWsHead(
}
override fun start() {
super.start()
this.subscription?.dispose()
val heads = Flux.merge(
// get the current block, not just wait for the next update
@@ -50,6 +51,7 @@ class EthereumWsHead(
}
override fun stop() {
super.stop()
subscription?.dispose()
subscription = null
}

View File

@@ -37,8 +37,6 @@ import org.springframework.context.Lifecycle
import org.springframework.util.ConcurrentReferenceHashMap
import reactor.core.publisher.Flux
import reactor.core.publisher.Mono
import java.util.concurrent.Executors
import java.util.concurrent.TimeUnit
import java.util.concurrent.locks.ReentrantLock
@Suppress("UNCHECKED_CAST")
@@ -70,20 +68,6 @@ open class EthereumPosMultiStream(
head = updateHead()
}
super.init()
Executors.newScheduledThreadPool(1).scheduleAtFixedRate({
val timeout = System.currentTimeMillis() - (head?.getLastUpdateTime() ?: 0L)
log.debug("Check head is active! Lst updated $timeout ms ago")
if (timeout > 60_000 && lock.tryLock()) {
log.warn("Timeout is over 1 min - restart head")
try {
head = updateHead()
} catch (e: Exception) {
log.error(e.message, e)
} finally {
lock.unlock()
}
}
}, 60, 30, TimeUnit.SECONDS)
}
override fun start() {

View File

@@ -87,7 +87,7 @@ class GrpcHead(
log.warn("Disconnected $chain from ${parent.getId()}: ${err.message}")
parent.setStatus(UpstreamAvailability.UNAVAILABLE)
Mono.empty<BlockchainOuterClass.ChainHead>()
}
}.doFinally { log.warn("Head subscription finished: $it") }
}
/**
@@ -116,10 +116,12 @@ class GrpcHead(
}
override fun start() {
super.start()
this.internalStart(remote)
}
override fun stop() {
super.stop()
headSubscription?.dispose()
}
}

View File

@@ -69,7 +69,12 @@ class EthereumHeadMock implements Head {
}
@Override
long getLastUpdateTime() {
return 0
void start() {
}
@Override
void stop() {
}
}

View File

@@ -60,6 +60,7 @@ class AbstractHeadSpec extends Specification {
.expectNext(blocks[1])
.then {
assert called
head.stop()
source.tryEmitComplete()
}
.expectComplete()
@@ -83,7 +84,10 @@ class AbstractHeadSpec extends Specification {
.expectNext(blocks[2])
.then { source.tryEmitNext(blocks[3]) }
.expectNext(blocks[3])
.then { source.tryEmitComplete() }
.then {
head.stop()
source.tryEmitComplete()
}
.expectComplete()
.verify(Duration.ofSeconds(1))
}
@@ -110,7 +114,10 @@ class AbstractHeadSpec extends Specification {
.then { source.tryEmitNext(wrongblock) }
.then { source.tryEmitNext(blocks[3]) }
.expectNext(blocks[3])
.then { source.tryEmitComplete() }
.then {
head.stop()
source.tryEmitComplete()
}
.expectComplete()
.verify(Duration.ofSeconds(1))
}
@@ -132,7 +139,7 @@ class AbstractHeadSpec extends Specification {
BlockContainer getHead() {
return null
}
}, new BlockValidator.AlwaysValid())
}, new BlockValidator.AlwaysValid(), 100_000)
}
}
}

View File

@@ -46,14 +46,14 @@ class MergedHeadSpec extends Specification {
class TestHead1 extends AbstractHead {
TestHead1() {
super(new MostWorkForkChoice())
super(new MostWorkForkChoice(), new BlockValidator.AlwaysValid(), 100_000)
}
}
class TestHead2 extends AbstractHead implements Lifecycle {
TestHead2() {
super(new MostWorkForkChoice())
super(new MostWorkForkChoice(), new BlockValidator.AlwaysValid(), 100_000)
}
@Override