try fix tests
This commit is contained in:
@@ -20,6 +20,7 @@ import io.emeraldpay.dshackle.upstream.forkchoice.ForkChoice
|
|||||||
import org.slf4j.LoggerFactory
|
import org.slf4j.LoggerFactory
|
||||||
import reactor.core.Disposable
|
import reactor.core.Disposable
|
||||||
import reactor.core.publisher.Flux
|
import reactor.core.publisher.Flux
|
||||||
|
import reactor.core.publisher.SignalType
|
||||||
import reactor.core.publisher.Sinks
|
import reactor.core.publisher.Sinks
|
||||||
import reactor.core.publisher.Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER
|
import reactor.core.publisher.Sinks.EmitResult.FAIL_ZERO_SUBSCRIBER
|
||||||
import reactor.core.publisher.Sinks.EmitResult.OK
|
import reactor.core.publisher.Sinks.EmitResult.OK
|
||||||
@@ -27,6 +28,7 @@ import reactor.core.scheduler.Schedulers
|
|||||||
import reactor.kotlin.core.publisher.toMono
|
import reactor.kotlin.core.publisher.toMono
|
||||||
import java.util.concurrent.Executors
|
import java.util.concurrent.Executors
|
||||||
import java.util.concurrent.TimeUnit
|
import java.util.concurrent.TimeUnit
|
||||||
|
import java.util.concurrent.locks.ReentrantLock
|
||||||
|
|
||||||
abstract class AbstractHead(
|
abstract class AbstractHead(
|
||||||
private val forkChoice: ForkChoice,
|
private val forkChoice: ForkChoice,
|
||||||
@@ -43,6 +45,7 @@ abstract class AbstractHead(
|
|||||||
private val beforeBlockHandlers = ArrayList<Runnable>()
|
private val beforeBlockHandlers = ArrayList<Runnable>()
|
||||||
private var stopping = false
|
private var stopping = false
|
||||||
private var lastHeadUpdated = 0L
|
private var lastHeadUpdated = 0L
|
||||||
|
private val lock = ReentrantLock()
|
||||||
|
|
||||||
init {
|
init {
|
||||||
Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(
|
Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(
|
||||||
@@ -50,7 +53,12 @@ abstract class AbstractHead(
|
|||||||
val delay = System.currentTimeMillis() - lastHeadUpdated
|
val delay = System.currentTimeMillis() - lastHeadUpdated
|
||||||
if (delay > awaitHeadTimeoutMs) {
|
if (delay > awaitHeadTimeoutMs) {
|
||||||
log.warn("No head updates for $delay ms @ ${this.javaClass} - restart")
|
log.warn("No head updates for $delay ms @ ${this.javaClass} - restart")
|
||||||
start()
|
try {
|
||||||
|
lock.tryLock()
|
||||||
|
start()
|
||||||
|
} finally {
|
||||||
|
lock.unlock()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}, 300, 30, TimeUnit.SECONDS
|
}, 300, 30, TimeUnit.SECONDS
|
||||||
)
|
)
|
||||||
@@ -71,13 +79,13 @@ abstract class AbstractHead(
|
|||||||
// close internal stream if upstream is finished, otherwise it gets stuck,
|
// close internal stream if upstream is finished, otherwise it gets stuck,
|
||||||
// but technically it should never happen during normal work, only when the Head
|
// but technically it should never happen during normal work, only when the Head
|
||||||
// is stopping
|
// is stopping
|
||||||
if (stopping) {
|
if (it == SignalType.ON_ERROR && !stopping) {
|
||||||
log.info("Received signal $it - stop emit new head!!!")
|
log.warn("Received signal $it unexpectedly - restart head")
|
||||||
|
lastHeadUpdated = 0L
|
||||||
|
} else {
|
||||||
|
log.warn("Received signal $it - stop emit new head!!!")
|
||||||
completed = true
|
completed = true
|
||||||
stream.tryEmitComplete()
|
stream.tryEmitComplete()
|
||||||
} else {
|
|
||||||
log.warn("Received signal $it unexpectedly - restart head")
|
|
||||||
start()
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
.subscribeOn(Schedulers.boundedElastic())
|
.subscribeOn(Schedulers.boundedElastic())
|
||||||
|
|||||||
Reference in New Issue
Block a user